[infinispan-commits] Infinispan SVN: r267 - in trunk: core/src/main/java/org/infinispan/atomic/atomichashmap and 20 other directories.
infinispan-commits at lists.jboss.org
infinispan-commits at lists.jboss.org
Tue May 12 23:15:11 EDT 2009
Author: mircea.markus
Date: 2009-05-12 23:15:11 -0400 (Tue, 12 May 2009)
New Revision: 267
Added:
trunk/core/src/main/java/org/infinispan/context/InvocationContextContainer.java
trunk/core/src/main/java/org/infinispan/context/InvocationContextContainerImpl.java
trunk/core/src/main/java/org/infinispan/context/impl/LocalTxInvocationContext.java
trunk/core/src/main/java/org/infinispan/transaction/xa/CacheTransaction.java
trunk/core/src/main/java/org/infinispan/transaction/xa/RemoteTransaction.java
trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionTable.java
trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionXaAdapter.java
Removed:
trunk/core/src/main/java/org/infinispan/context/container/
trunk/core/src/main/java/org/infinispan/context/impl/InitiatorTxInvocationContext.java
trunk/core/src/main/java/org/infinispan/loader/
trunk/core/src/main/java/org/infinispan/lock/
trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionXaAdapter.java
trunk/core/src/main/java/org/infinispan/transaction/xa/TxEnlistingManager.java
Modified:
trunk/core/src/main/java/org/infinispan/AbstractDelegatingAdvancedCache.java
trunk/core/src/main/java/org/infinispan/AdvancedCache.java
trunk/core/src/main/java/org/infinispan/CacheDelegate.java
trunk/core/src/main/java/org/infinispan/atomic/atomichashmap/AtomicHashMap.java
trunk/core/src/main/java/org/infinispan/atomic/atomichashmap/AtomicHashMapProxy.java
trunk/core/src/main/java/org/infinispan/commands/CommandsFactoryImpl.java
trunk/core/src/main/java/org/infinispan/commands/remote/BaseRpcCommand.java
trunk/core/src/main/java/org/infinispan/commands/remote/ClusteredGetCommand.java
trunk/core/src/main/java/org/infinispan/commands/tx/AbstractTransactionBoundaryCommand.java
trunk/core/src/main/java/org/infinispan/commands/tx/PrepareCommand.java
trunk/core/src/main/java/org/infinispan/context/impl/RemoteTxInvocationContext.java
trunk/core/src/main/java/org/infinispan/context/impl/TxInvocationContext.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/interceptors/BatchingInterceptor.java
trunk/core/src/main/java/org/infinispan/interceptors/InterceptorChain.java
trunk/core/src/main/java/org/infinispan/interceptors/NotificationInterceptor.java
trunk/core/src/main/java/org/infinispan/interceptors/TxInterceptor.java
trunk/core/src/main/java/org/infinispan/notifications/cachelistener/CacheNotifierImpl.java
trunk/core/src/main/java/org/infinispan/remoting/InboundInvocationHandlerImpl.java
trunk/core/src/main/java/org/infinispan/statetransfer/StateTransferManagerImpl.java
trunk/core/src/main/java/org/infinispan/util/concurrent/locks/LockManagerImpl.java
trunk/core/src/main/java/org/infinispan/util/concurrent/locks/OwnableReentrantLock.java
trunk/core/src/main/java/org/infinispan/util/concurrent/locks/containers/OwnableReentrantPerEntryLockContainer.java
trunk/core/src/main/java/org/infinispan/util/concurrent/locks/containers/OwnableReentrantStripedLockContainer.java
trunk/core/src/test/java/org/infinispan/api/mvcc/LockAssert.java
trunk/core/src/test/java/org/infinispan/api/mvcc/LockPerEntryTest.java
trunk/core/src/test/java/org/infinispan/api/mvcc/LockTestBase.java
trunk/core/src/test/java/org/infinispan/api/mvcc/PutForExternalReadTest.java
trunk/core/src/test/java/org/infinispan/api/mvcc/repeatable_read/WriteSkewTest.java
trunk/core/src/test/java/org/infinispan/atomic/AtomicMapFunctionalTest.java
trunk/core/src/test/java/org/infinispan/notifications/cachelistener/CacheNotifierImplTest.java
trunk/core/src/test/java/org/infinispan/test/AbstractCacheTest.java
trunk/core/src/test/java/org/infinispan/test/TestingUtil.java
trunk/tree/src/main/java/org/infinispan/tree/NodeImpl.java
trunk/tree/src/main/java/org/infinispan/tree/TreeCacheImpl.java
trunk/tree/src/main/java/org/infinispan/tree/TreeStructureSupport.java
trunk/tree/src/test/java/org/infinispan/api/tree/NodeMoveAPITest.java
Log:
fixed UT and refined code after code review
Modified: trunk/core/src/main/java/org/infinispan/AbstractDelegatingAdvancedCache.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/AbstractDelegatingAdvancedCache.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/AbstractDelegatingAdvancedCache.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -3,7 +3,7 @@
import org.infinispan.batch.BatchContainer;
import org.infinispan.container.DataContainer;
import org.infinispan.context.Flag;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.eviction.EvictionManager;
import org.infinispan.factories.ComponentRegistry;
import org.infinispan.interceptors.base.CommandInterceptor;
Modified: trunk/core/src/main/java/org/infinispan/AdvancedCache.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/AdvancedCache.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/AdvancedCache.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -3,7 +3,7 @@
import org.infinispan.batch.BatchContainer;
import org.infinispan.container.DataContainer;
import org.infinispan.context.Flag;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.eviction.EvictionManager;
import org.infinispan.factories.ComponentRegistry;
import org.infinispan.interceptors.base.CommandInterceptor;
Modified: trunk/core/src/main/java/org/infinispan/CacheDelegate.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/CacheDelegate.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/CacheDelegate.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -41,7 +41,7 @@
import org.infinispan.container.entries.InternalCacheEntry;
import org.infinispan.context.Flag;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.eviction.EvictionManager;
import org.infinispan.factories.ComponentRegistry;
import org.infinispan.factories.annotations.Inject;
@@ -139,8 +139,7 @@
}
public final boolean remove(Object key, Object value) {
- RemoveCommand command = commandsFactory.buildRemoveCommand(key, value);
- return (Boolean) invoker.invoke(getInvocationContext(), command);
+ return remove(key, value, (Flag[]) null);
}
public final boolean replace(K key, V oldValue, V newValue) {
@@ -161,9 +160,7 @@
}
public final boolean containsKey(Object key) {
- GetKeyValueCommand command = commandsFactory.buildGetKeyValueCommand(key);
- Object response = invoker.invoke(getInvocationContext(), command);
- return response != null;
+ return containsKey(key, (Flag[]) null);
}
public final boolean containsValue(Object value) {
@@ -172,8 +169,7 @@
@SuppressWarnings("unchecked")
public final V get(Object key) {
- GetKeyValueCommand command = commandsFactory.buildGetKeyValueCommand(key);
- return (V) invoker.invoke(getInvocationContext(), command);
+ return get(key, (Flag[])null);
}
public final V put(K key, V value) {
@@ -182,8 +178,7 @@
@SuppressWarnings("unchecked")
public final V remove(Object key) {
- RemoveCommand command = commandsFactory.buildRemoveCommand(key, null);
- return (V) invoker.invoke(getInvocationContext(), command);
+ return remove(key, (Flag[]) null);
}
public final void putAll(Map<? extends K, ? extends V> map) {
@@ -191,8 +186,7 @@
}
public final void clear() {
- ClearCommand command = commandsFactory.buildClearCommand();
- invoker.invoke(getInvocationContext(), command);
+ clear((Flag[]) null);
}
public Set<K> keySet() {
@@ -233,7 +227,7 @@
}
private InvocationContext getInvocationContext() {
- return icc.getLocalInvocationContext(true);
+ return icc.getLocalInvocationContext();
}
public void lock(K key) {
@@ -333,8 +327,10 @@
}
public final V put(K key, V value, long lifespan, TimeUnit lifespanUnit, long maxIdleTime, TimeUnit maxIdleTimeUnit, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return put(key, value, lifespan, lifespanUnit, maxIdleTime, maxIdleTimeUnit);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ PutKeyValueCommand command = commandsFactory.buildPutKeyValueCommand(key, value, lifespanUnit.toMillis(lifespan), maxIdleTimeUnit.toMillis(maxIdleTime));
+ return (V) invoker.invoke(ctx, command);
}
public final V putIfAbsent(K key, V value, Flag... flags) {
@@ -342,8 +338,11 @@
}
public final V putIfAbsent(K key, V value, long lifespan, TimeUnit lifespanUnit, long maxIdleTime, TimeUnit maxIdleTimeUnit, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return putIfAbsent(key, value, lifespan, lifespanUnit, maxIdleTime, maxIdleTimeUnit);
+ InvocationContext context = getInvocationContext();
+ context.setFlags(flags);
+ PutKeyValueCommand command = commandsFactory.buildPutKeyValueCommand(key, value, lifespanUnit.toMillis(lifespan), maxIdleTimeUnit.toMillis(maxIdleTime));
+ command.setPutIfAbsent(true);
+ return (V) invoker.invoke(context, command);
}
public final void putAll(Map<? extends K, ? extends V> map, Flag... flags) {
@@ -351,23 +350,31 @@
}
public final void putAll(Map<? extends K, ? extends V> map, long lifespan, TimeUnit lifespanUnit, long maxIdleTime, TimeUnit maxIdleTimeUnit, Flag... flags) {
- getInvocationContext().setFlags(flags);
- putAll(map);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ PutMapCommand command = commandsFactory.buildPutMapCommand(map, MILLISECONDS.toMillis(defaultLifespan), MILLISECONDS.toMillis(defaultMaxIdleTime));
+ invoker.invoke(ctx, command);
}
public final V remove(Object key, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return remove(key);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ RemoveCommand command = commandsFactory.buildRemoveCommand(key, null);
+ return (V) invoker.invoke(ctx, command);
}
public final boolean remove(Object key, Object oldValue, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return remove(key, oldValue);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ RemoveCommand command = commandsFactory.buildRemoveCommand(key, oldValue);
+ return (Boolean) invoker.invoke(ctx, command);
}
public final void clear(Flag... flags) {
- getInvocationContext().setFlags(flags);
- clear();
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ ClearCommand command = commandsFactory.buildClearCommand();
+ invoker.invoke(ctx, command);
}
public final V replace(K k, V v, Flag... flags) {
@@ -379,13 +386,17 @@
}
public final V replace(K k, V v, long lifespan, TimeUnit lifespanUnit, long maxIdle, TimeUnit maxIdleUnit, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return replace(k, v, lifespan, lifespanUnit, maxIdle, maxIdleUnit);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ ReplaceCommand command = commandsFactory.buildReplaceCommand(k, null, v, lifespanUnit.toMillis(lifespan), maxIdleUnit.toMillis(maxIdle));
+ return (V) invoker.invoke(ctx, command);
}
public final boolean replace(K k, V oV, V nV, long lifespan, TimeUnit lifespanUnit, long maxIdle, TimeUnit maxIdleUnit, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return replace(k, oV, nV, lifespan, lifespanUnit, maxIdle, maxIdleUnit);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ ReplaceCommand command = commandsFactory.buildReplaceCommand(k, oV, nV, lifespanUnit.toMillis(lifespan), maxIdleUnit.toMillis(maxIdle));
+ return (Boolean) invoker.invoke(ctx, command);
}
public final Future<V> putAsync(K key, V value, Flag... flags) {
@@ -393,8 +404,11 @@
}
public final Future<V> putAsync(K key, V value, long lifespan, TimeUnit lifespanUnit, long maxIdleTime, TimeUnit maxIdleTimeUnit, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return putAsync(key, value, lifespan, lifespanUnit, maxIdleTime, maxIdleTimeUnit);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ ctx.setUseFutureReturnType(true);
+ PutKeyValueCommand command = commandsFactory.buildPutKeyValueCommand(key, value, lifespanUnit.toMillis(lifespan), maxIdleTimeUnit.toMillis(maxIdleTime));
+ return wrapInFuture(invoker.invoke(ctx, command));
}
public final Future<V> putIfAbsentAsync(K key, V value, Flag... flags) {
@@ -402,8 +416,12 @@
}
public final Future<V> putIfAbsentAsync(K key, V value, long lifespan, TimeUnit lifespanUnit, long maxIdleTime, TimeUnit maxIdleTimeUnit, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return putIfAbsentAsync(key, value, lifespan, lifespanUnit, maxIdleTime, maxIdleTimeUnit);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ ctx.setUseFutureReturnType(true);
+ PutKeyValueCommand command = commandsFactory.buildPutKeyValueCommand(key, value, lifespanUnit.toMillis(lifespan), maxIdleTimeUnit.toMillis(maxIdleTime));
+ command.setPutIfAbsent(true);
+ return wrapInFuture(invoker.invoke(ctx, command));
}
public final Future<Void> putAllAsync(Map<? extends K, ? extends V> map, Flag... flags) {
@@ -411,18 +429,27 @@
}
public final Future<Void> putAllAsync(Map<? extends K, ? extends V> map, long lifespan, TimeUnit lifespanUnit, long maxIdleTime, TimeUnit maxIdleTimeUnit, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return putAllAsync(map, MILLISECONDS.toMillis(defaultLifespan), MILLISECONDS, MILLISECONDS.toMillis(defaultMaxIdleTime), MILLISECONDS);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ ctx.setUseFutureReturnType(true);
+ PutMapCommand command = commandsFactory.buildPutMapCommand(map, MILLISECONDS.toMillis(MILLISECONDS.toMillis(defaultLifespan)), MILLISECONDS.toMillis(MILLISECONDS.toMillis(defaultMaxIdleTime)));
+ return wrapInFuture(invoker.invoke(ctx, command));
}
public final Future<V> removeAsync(Object key, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return removeAsync(key);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ ctx.setUseFutureReturnType(true);
+ RemoveCommand command = commandsFactory.buildRemoveCommand(key, null);
+ return wrapInFuture(invoker.invoke(ctx, command));
}
public final Future<Void> clearAsync(Flag... flags) {
- getInvocationContext().setFlags(flags);
- return clearAsync();
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ ctx.setUseFutureReturnType(true);
+ ClearCommand command = commandsFactory.buildClearCommand();
+ return wrapInFuture(invoker.invoke(ctx, command));
}
public final Future<V> replaceAsync(K k, V v, Flag... flags) {
@@ -434,23 +461,34 @@
}
public final Future<V> replaceAsync(K k, V v, long lifespan, TimeUnit lifespanUnit, long maxIdle, TimeUnit maxIdleUnit, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return replaceAsync(k, v, lifespan, lifespanUnit, maxIdle, maxIdleUnit);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ ctx.setUseFutureReturnType(true);
+ ReplaceCommand command = commandsFactory.buildReplaceCommand(k, null, v, lifespanUnit.toMillis(lifespan), maxIdleUnit.toMillis(maxIdle));
+ return wrapInFuture(invoker.invoke(ctx, command));
}
public final Future<Boolean> replaceAsync(K k, V oV, V nV, long lifespan, TimeUnit lifespanUnit, long maxIdle, TimeUnit maxIdleUnit, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return replaceAsync(k, oV, nV, lifespan, lifespanUnit, maxIdle, maxIdleUnit);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ ctx.setUseFutureReturnType(true);
+ ReplaceCommand command = commandsFactory.buildReplaceCommand(k, oV, nV, lifespanUnit.toMillis(lifespan), maxIdleUnit.toMillis(maxIdle));
+ return wrapInFuture(invoker.invoke(ctx, command));
}
public final boolean containsKey(Object key, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return containsKey(key);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ GetKeyValueCommand command = commandsFactory.buildGetKeyValueCommand(key);
+ Object response = invoker.invoke(ctx, command);
+ return response != null;
}
public final V get(Object key, Flag... flags) {
- getInvocationContext().setFlags(flags);
- return get(key);
+ InvocationContext ctx = getInvocationContext();
+ ctx.setFlags(flags);
+ GetKeyValueCommand command = commandsFactory.buildGetKeyValueCommand(key);
+ return (V) invoker.invoke(ctx, command);
}
public ComponentStatus getStatus() {
@@ -507,15 +545,12 @@
@SuppressWarnings("unchecked")
public final V put(K key, V value, long lifespan, TimeUnit lifespanUnit, long maxIdleTime, TimeUnit idleTimeUnit) {
- PutKeyValueCommand command = commandsFactory.buildPutKeyValueCommand(key, value, lifespanUnit.toMillis(lifespan), idleTimeUnit.toMillis(maxIdleTime));
- return (V) invoker.invoke(getInvocationContext(), command);
+ return put(key, value, lifespan, lifespanUnit, maxIdleTime, idleTimeUnit, (Flag[]) null);
}
@SuppressWarnings("unchecked")
public final V putIfAbsent(K key, V value, long lifespan, TimeUnit lifespanUnit, long maxIdleTime, TimeUnit idleTimeUnit) {
- PutKeyValueCommand command = commandsFactory.buildPutKeyValueCommand(key, value, lifespanUnit.toMillis(lifespan), idleTimeUnit.toMillis(maxIdleTime));
- command.setPutIfAbsent(true);
- return (V) invoker.invoke(getInvocationContext(), command);
+ return putIfAbsent(key, value, lifespan, lifespanUnit, maxIdleTime, idleTimeUnit, (Flag[]) null);
}
public final void putAll(Map<? extends K, ? extends V> map, long lifespan, TimeUnit lifespanUnit, long maxIdleTime, TimeUnit idleTimeUnit) {
@@ -525,8 +560,7 @@
@SuppressWarnings("unchecked")
public final V replace(K key, V value, long lifespan, TimeUnit lifespanUnit, long maxIdleTime, TimeUnit idleTimeUnit) {
- ReplaceCommand command = commandsFactory.buildReplaceCommand(key, null, value, lifespanUnit.toMillis(lifespan), idleTimeUnit.toMillis(maxIdleTime));
- return (V) invoker.invoke(getInvocationContext(), command);
+ return replace(key, value, lifespan, lifespanUnit, maxIdleTime, idleTimeUnit, (Flag[]) null);
}
public final boolean replace(K key, V oldValue, V value, long lifespan, TimeUnit lifespanUnit, long maxIdleTime, TimeUnit idleTimeUnit) {
@@ -652,10 +686,7 @@
}
public final Future<V> replaceAsync(K key, V value, long lifespan, TimeUnit lifespanUnit, long maxIdle, TimeUnit maxIdleUnit) {
- InvocationContext ctx = getInvocationContext();
- ctx.setUseFutureReturnType(true);
- ReplaceCommand command = commandsFactory.buildReplaceCommand(key, null, value, lifespanUnit.toMillis(lifespan), maxIdleUnit.toMillis(maxIdle));
- return wrapInFuture(invoker.invoke(getInvocationContext(), command));
+ return replaceAsync(key, value, lifespan, lifespanUnit, maxIdle, maxIdleUnit, (Flag[]) null);
}
public final Future<Boolean> replaceAsync(K key, V oldValue, V newValue) {
Modified: trunk/core/src/main/java/org/infinispan/atomic/atomichashmap/AtomicHashMap.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/atomic/atomichashmap/AtomicHashMap.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/atomic/atomichashmap/AtomicHashMap.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -31,7 +31,7 @@
import org.infinispan.atomic.Delta;
import org.infinispan.atomic.NullDelta;
import org.infinispan.batch.BatchContainer;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.util.FastCopyHashMap;
import java.util.Collection;
Modified: trunk/core/src/main/java/org/infinispan/atomic/atomichashmap/AtomicHashMapProxy.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/atomic/atomichashmap/AtomicHashMapProxy.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/atomic/atomichashmap/AtomicHashMapProxy.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -27,7 +27,7 @@
import org.infinispan.batch.BatchContainer;
import org.infinispan.context.Flag;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.util.logging.Log;
import org.infinispan.util.logging.LogFactory;
@@ -125,7 +125,7 @@
public V put(K key, V value) {
try {
startAtomic();
- InvocationContext ctx = icc.getLocalInvocationContext(true);
+ InvocationContext ctx = icc.getLocalInvocationContext();
AtomicHashMap<K, V> deltaMapForWrite = getDeltaMapForWrite(ctx);
return deltaMapForWrite.put(key, value);
}
@@ -137,7 +137,7 @@
public V remove(Object key) {
try {
startAtomic();
- InvocationContext ic = icc.getLocalInvocationContext(true);
+ InvocationContext ic = icc.getLocalInvocationContext();
return getDeltaMapForWrite(ic).remove(key);
}
finally {
@@ -148,7 +148,7 @@
public void putAll(Map<? extends K, ? extends V> m) {
try {
startAtomic();
- InvocationContext ic = icc.getLocalInvocationContext(true);
+ InvocationContext ic = icc.getLocalInvocationContext();
getDeltaMapForWrite(ic).putAll(m);
}
finally {
@@ -159,7 +159,7 @@
public void clear() {
try {
startAtomic();
- InvocationContext ic = icc.getLocalInvocationContext(true);
+ InvocationContext ic = icc.getLocalInvocationContext();
getDeltaMapForWrite(ic).clear();
}
finally {
Modified: trunk/core/src/main/java/org/infinispan/commands/CommandsFactoryImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/commands/CommandsFactoryImpl.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/commands/CommandsFactoryImpl.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -41,7 +41,7 @@
import org.infinispan.commands.write.ReplaceCommand;
import org.infinispan.commands.write.WriteCommand;
import org.infinispan.container.DataContainer;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.distribution.DistributionManager;
import org.infinispan.factories.annotations.Inject;
import org.infinispan.factories.annotations.Start;
@@ -49,6 +49,7 @@
import org.infinispan.loaders.CacheLoaderManager;
import org.infinispan.notifications.cachelistener.CacheNotifier;
import org.infinispan.transaction.xa.GlobalTransaction;
+import org.infinispan.transaction.xa.TransactionTable;
import java.util.Collection;
import java.util.List;
@@ -70,11 +71,13 @@
private InterceptorChain interceptorChain;
private DistributionManager distributionManager;
private InvocationContextContainer icc;
+ private TransactionTable txTable;
@Inject
public void setupDependencies(DataContainer container, CacheNotifier notifier, Cache cache,
InterceptorChain interceptorChain, CacheLoaderManager clManager,
- DistributionManager distributionManager, InvocationContextContainer icc) {
+ DistributionManager distributionManager, InvocationContextContainer icc,
+ TransactionTable txTable) {
this.dataContainer = container;
this.notifier = notifier;
this.cache = cache;
@@ -82,6 +85,7 @@
this.cacheLoaderManager = clManager;
this.distributionManager = distributionManager;
this.icc = icc;
+ this.txTable = txTable;
}
@Start(priority = 1)
@@ -205,18 +209,18 @@
break;
case PrepareCommand.COMMAND_ID:
PrepareCommand pc = (PrepareCommand) c;
- pc.init(interceptorChain, icc);
+ pc.init(interceptorChain, icc, txTable);
pc.initialize(notifier);
if (pc.getModifications() != null)
for (ReplicableCommand nested : pc.getModifications()) initializeReplicableCommand(nested);
break;
case CommitCommand.COMMAND_ID:
CommitCommand commitCommand = (CommitCommand) c;
- commitCommand.init(interceptorChain, icc);
+ commitCommand.init(interceptorChain, icc, txTable);
break;
case RollbackCommand.COMMAND_ID:
RollbackCommand rollbackCommand = (RollbackCommand) c;
- rollbackCommand.init(interceptorChain, icc);
+ rollbackCommand.init(interceptorChain, icc, txTable);
break;
case ClearCommand.COMMAND_ID:
ClearCommand cc = (ClearCommand) c;
Modified: trunk/core/src/main/java/org/infinispan/commands/remote/BaseRpcCommand.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/commands/remote/BaseRpcCommand.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/commands/remote/BaseRpcCommand.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -6,7 +6,7 @@
import org.infinispan.util.logging.Log;
import org.infinispan.util.logging.LogFactory;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
/**
* Base class for RPC commands.
Modified: trunk/core/src/main/java/org/infinispan/commands/remote/ClusteredGetCommand.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/commands/remote/ClusteredGetCommand.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/commands/remote/ClusteredGetCommand.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -26,7 +26,7 @@
import org.infinispan.container.entries.CacheEntry;
import org.infinispan.container.entries.InternalCacheEntry;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.loaders.CacheLoaderManager;
import org.infinispan.util.logging.Log;
import org.infinispan.util.logging.LogFactory;
Modified: trunk/core/src/main/java/org/infinispan/commands/tx/AbstractTransactionBoundaryCommand.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/commands/tx/AbstractTransactionBoundaryCommand.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/commands/tx/AbstractTransactionBoundaryCommand.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -22,10 +22,12 @@
package org.infinispan.commands.tx;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.context.impl.RemoteTxInvocationContext;
import org.infinispan.interceptors.InterceptorChain;
import org.infinispan.transaction.xa.GlobalTransaction;
+import org.infinispan.transaction.xa.TransactionTable;
+import org.infinispan.transaction.xa.RemoteTransaction;
import org.infinispan.util.logging.Log;
import org.infinispan.util.logging.LogFactory;
@@ -44,10 +46,13 @@
protected String cacheName;
protected InterceptorChain invoker;
protected InvocationContextContainer icc;
+ protected TransactionTable txTable;
- public void init(InterceptorChain chain, InvocationContextContainer icc) {
+
+ public void init(InterceptorChain chain, InvocationContextContainer icc, TransactionTable txTable) {
this.invoker = chain;
this.icc = icc;
+ this.txTable = txTable;
}
public String getCacheName() {
@@ -68,12 +73,18 @@
public Object perform(InvocationContext ctx) throws Throwable {
if (ctx != null) throw new IllegalStateException("Expected null context!");
- RemoteTxInvocationContext ctxt = icc.getRemoteTxInvocationContext(getGlobalTransaction(), false);
- if (ctxt == null) {
- if (log.isInfoEnabled()) log.info("Not found RemoteTxInvocationContext for tx: " + getGlobalTransaction());
+ RemoteTransaction transaction = txTable.getRemoteTransaction(globalTx);
+ if (transaction == null) {
+ if (log.isInfoEnabled()) log.info("Not found RemoteTransaction for tx id: " + globalTx);
return null;
}
- return invoker.invoke(ctxt, this);
+ RemoteTxInvocationContext ctxt = icc.getRemoteTxInvocationContext();
+ ctxt.setRemoteTransaction(transaction);
+ try {
+ return invoker.invoke(ctxt, this);
+ } finally {
+ txTable.removeRemoteTransaction(globalTx);
+ }
}
public Object[] getParameters() {
Modified: trunk/core/src/main/java/org/infinispan/commands/tx/PrepareCommand.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/commands/tx/PrepareCommand.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/commands/tx/PrepareCommand.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -29,6 +29,7 @@
import org.infinispan.context.impl.TxInvocationContext;
import org.infinispan.notifications.cachelistener.CacheNotifier;
import org.infinispan.transaction.xa.GlobalTransaction;
+import org.infinispan.transaction.xa.RemoteTransaction;
import java.util.Arrays;
import java.util.List;
@@ -69,9 +70,15 @@
public final Object perform(InvocationContext ignored) throws Throwable {
if (ignored != null) throw new IllegalStateException("Expected null context!");
- RemoteTxInvocationContext ctx = icc.getRemoteTxInvocationContext(getGlobalTransaction(), true);
- notifier.notifyTransactionRegistered(ctx.getClusterTransactionId(), ctx);
- ctx.initialize(modifications, globalTx);
+
+ //1. first create a remote transaction
+ RemoteTransaction remoteTransaction = txTable.createRemoteTransaction(globalTx, modifications);
+
+ //2. then set it on the invocation context
+ RemoteTxInvocationContext ctx = icc.getRemoteTxInvocationContext();
+ ctx.setRemoteTransaction(remoteTransaction);
+
+ notifier.notifyTransactionRegistered(ctx.getGlobalTransaction(), ctx);
return invoker.invoke(ctx, this);
}
Copied: trunk/core/src/main/java/org/infinispan/context/InvocationContextContainer.java (from rev 252, trunk/core/src/main/java/org/infinispan/context/container/InvocationContextContainer.java)
===================================================================
--- trunk/core/src/main/java/org/infinispan/context/InvocationContextContainer.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/context/InvocationContextContainer.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -0,0 +1,32 @@
+package org.infinispan.context;
+
+import org.infinispan.context.impl.LocalTxInvocationContext;
+import org.infinispan.context.impl.RemoteTxInvocationContext;
+import org.infinispan.factories.annotations.NonVolatile;
+import org.infinispan.factories.scopes.Scope;
+import org.infinispan.factories.scopes.Scopes;
+
+/**
+ * // TODO: Mircea: Document this!
+ *
+ * @author Manik Surtani (<a href="mailto:manik at jboss.org">manik at jboss.org</a>)
+ * @author Mircea.Markus at jboss.com
+ * @since 4.0
+ */
+ at NonVolatile
+ at Scope(Scopes.NAMED_CACHE)
+public interface InvocationContextContainer {
+ InvocationContext getLocalInvocationContext();
+
+ LocalTxInvocationContext getInitiatorTxInvocationContext();
+
+ RemoteTxInvocationContext getRemoteTxInvocationContext();
+
+ InvocationContext getRemoteNonTxInvocationContext();
+
+ InvocationContext getThreadContext();
+
+ Object suspend();
+
+ void resume(Object backup);
+}
Property changes on: trunk/core/src/main/java/org/infinispan/context/InvocationContextContainer.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Copied: trunk/core/src/main/java/org/infinispan/context/InvocationContextContainerImpl.java (from rev 252, trunk/core/src/main/java/org/infinispan/context/container/ReplicationInvocationContextContainer.java)
===================================================================
--- trunk/core/src/main/java/org/infinispan/context/InvocationContextContainerImpl.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/context/InvocationContextContainerImpl.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -0,0 +1,171 @@
+/*
+ * JBoss, Home of Professional Open Source.
+ * Copyright 2000 - 2008, Red Hat Middleware LLC, and individual contributors
+ * as indicated by the @author tags. See the copyright.txt file in the
+ * distribution for a full listing of individual contributors.
+ *
+ * This is free software; you can redistribute it and/or modify it
+ * under the terms of the GNU Lesser General Public License as
+ * published by the Free Software Foundation; either version 2.1 of
+ * the License, or (at your option) any later version.
+ *
+ * This software is distributed in the hope that it will be useful,
+ * but WITHOUT ANY WARRANTY; without even the implied warranty of
+ * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the GNU
+ * Lesser General Public License for more details.
+ *
+ * You should have received a copy of the GNU Lesser General Public
+ * License along with this software; if not, write to the Free
+ * Software Foundation, Inc., 51 Franklin St, Fifth Floor, Boston, MA
+ * 02110-1301 USA, or see the FSF site: http://www.fsf.org.
+ */
+package org.infinispan.context;
+
+import org.infinispan.CacheException;
+import org.infinispan.context.impl.LocalTxInvocationContext;
+import org.infinispan.context.impl.NonTxInvocationContext;
+import org.infinispan.context.impl.RemoteTxInvocationContext;
+import org.infinispan.factories.annotations.Inject;
+import org.infinispan.transaction.xa.TransactionXaAdapter;
+import org.infinispan.transaction.xa.GlobalTransaction;
+import org.infinispan.transaction.xa.TransactionTable;
+import org.infinispan.util.logging.Log;
+import org.infinispan.util.logging.LogFactory;
+
+import javax.transaction.SystemException;
+import javax.transaction.Transaction;
+import javax.transaction.TransactionManager;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+
+
+/**
+ * Container and factory for thread locals
+ *
+ * @author Manik Surtani (<a href="mailto:manik at jboss.org">manik at jboss.org</a>)
+ * @since 4.0
+ */
+public class InvocationContextContainerImpl implements InvocationContextContainer {
+
+ private static Log log = LogFactory.getLog(InvocationContextContainerImpl.class);
+
+ private TransactionManager tm;
+ private TransactionTable transactionTable;
+
+
+ private ThreadLocal<PossibleContexts> contextsTl = new ThreadLocal<PossibleContexts>() {
+ @Override
+ protected PossibleContexts initialValue() {
+ return new PossibleContexts();
+ }
+ };
+
+ private Map<GlobalTransaction, RemoteTxInvocationContext> remoteTxMap = new ConcurrentHashMap<GlobalTransaction, RemoteTxInvocationContext>(20);
+
+ @Inject
+ public void init(TransactionManager tm, TransactionTable transactionTable) {
+ this.tm = tm;
+ this.transactionTable = transactionTable;
+ }
+
+ public InvocationContext getLocalInvocationContext() {
+ PossibleContexts contexts = contextsTl.get();
+ Transaction tx = getRunningTx();
+ if (tx != null) {
+ contexts.initInitiatorInvicationContext();
+ LocalTxInvocationContext context = contexts.localTxInvocationContext;
+ TransactionXaAdapter xaAdapter = transactionTable.getXaCacheAdapter(tx);
+ context.setXaCache(xaAdapter);
+ return contexts.updateThreadContextAndReturn(context);
+ } else {
+ contexts.initNonTxInvocationContext();
+ contexts.nonTxInvocationContext.prepareForCall();
+ contexts.nonTxInvocationContext.setOriginLocal(true);
+ return contexts.updateThreadContextAndReturn(contexts.nonTxInvocationContext);
+ }
+ }
+
+ public LocalTxInvocationContext getInitiatorTxInvocationContext() {
+ PossibleContexts contexts = contextsTl.get();
+ contexts.initInitiatorInvicationContext();
+ contexts.updateThreadContextAndReturn(contexts.localTxInvocationContext);
+ return contexts.localTxInvocationContext;
+ }
+
+ public RemoteTxInvocationContext getRemoteTxInvocationContext() {
+ PossibleContexts contexts = contextsTl.get();
+ contexts.initRemoteTxInvocationContext();
+ contexts.updateThreadContextAndReturn(contexts.remoteTxContext);
+ return contexts.remoteTxContext;
+ }
+
+ public InvocationContext getRemoteNonTxInvocationContext() {
+ PossibleContexts contexts = contextsTl.get();
+ contexts.initRemoteNonTxInvocationContext();
+ contexts.remoteNonTxContext.prepareForCall();
+ return contexts.updateThreadContextAndReturn(contexts.remoteNonTxContext);
+ }
+
+ public InvocationContext getThreadContext() {
+ InvocationContext invocationContext = contextsTl.get().threadInvocationContex;
+ if (invocationContext == null)
+ throw new IllegalStateException("This method can only be called after associating the current thread with a context");
+ return invocationContext;
+ }
+
+
+ public Object suspend() {
+ PossibleContexts result = contextsTl.get();
+ contextsTl.remove();
+ return result;
+ }
+
+ public void resume(Object backup) {
+ contextsTl.set((PossibleContexts) backup);
+ }
+
+ public static class PossibleContexts {
+ private NonTxInvocationContext nonTxInvocationContext;
+ private LocalTxInvocationContext localTxInvocationContext;
+ private NonTxInvocationContext remoteNonTxContext;
+ private RemoteTxInvocationContext remoteTxContext;
+ private InvocationContext threadInvocationContex;
+
+ public void initInitiatorInvicationContext() {
+ if (localTxInvocationContext == null) {
+ localTxInvocationContext = new LocalTxInvocationContext();
+ }
+ }
+
+ public void initNonTxInvocationContext() {
+ if (nonTxInvocationContext == null) {
+ nonTxInvocationContext = new NonTxInvocationContext();
+ }
+ }
+
+ public void initRemoteNonTxInvocationContext() {
+ if (remoteNonTxContext == null) {
+ remoteNonTxContext = new NonTxInvocationContext();
+ }
+ }
+
+ public void initRemoteTxInvocationContext() {
+ if (remoteTxContext == null) {
+ remoteTxContext = new RemoteTxInvocationContext();
+ }
+ }
+
+ public InvocationContext updateThreadContextAndReturn(InvocationContext ic) {
+ threadInvocationContex = ic;
+ return ic;
+ }
+ }
+
+ private Transaction getRunningTx() {
+ try {
+ return tm == null ? null : tm.getTransaction();
+ } catch (SystemException e) {
+ throw new CacheException(e);
+ }
+ }
+}
\ No newline at end of file
Property changes on: trunk/core/src/main/java/org/infinispan/context/InvocationContextContainerImpl.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Deleted: trunk/core/src/main/java/org/infinispan/context/impl/InitiatorTxInvocationContext.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/context/impl/InitiatorTxInvocationContext.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/context/impl/InitiatorTxInvocationContext.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -1,74 +0,0 @@
-package org.infinispan.context.impl;
-
-import org.infinispan.commands.write.WriteCommand;
-import org.infinispan.container.entries.CacheEntry;
-import org.infinispan.transaction.xa.GlobalTransaction;
-import org.infinispan.transaction.xa.TransactionXaAdapter;
-import org.infinispan.util.BidirectionalMap;
-
-import javax.transaction.Transaction;
-import java.util.List;
-import java.util.Map;
-
-/**
- * // TODO: Mircea: Document this!
- *
- * @author Mircea.Markus at jboss.com
- * @since 4.0
- */
-public class InitiatorTxInvocationContext extends AbstractTxInvocationContext {
-
- private TransactionXaAdapter xaAdapter;
-
- public Transaction getRunningTransaction() {
- return xaAdapter.getTransaction();
- }
-
- public boolean isOriginLocal() {
- return true;
- }
-
- public boolean isInTxScope() {
- return true;
- }
-
- public Object getLockOwner() {
- return xaAdapter.getTransactionIdentifier();
- }
-
- public GlobalTransaction getClusterTransactionId() {
- return xaAdapter.getTransactionIdentifier();
- }
-
- public List<WriteCommand> getModifications() {
- return xaAdapter.getModifications();
- }
-
- public void setXaCache(TransactionXaAdapter xaAdapter) {
- this.xaAdapter = xaAdapter;
- }
-
- public CacheEntry lookupEntry(Object key) {
- return xaAdapter != null ? xaAdapter.lookupEntry(key) : null;
- }
-
- public BidirectionalMap<Object, CacheEntry> getLookedUpEntries() {
- return xaAdapter.getLookedUpEntries();
- }
-
- public void putLookedUpEntry(Object key, CacheEntry e) {
- xaAdapter.putLookedUpEntry(key, e);
- }
-
- public void putLookedUpEntries(Map<Object, CacheEntry> lookedUpEntries) {
- xaAdapter.putLookedUpEntries(lookedUpEntries);
- }
-
- public void removeLookedUpEntry(Object key) {
- xaAdapter.removeLookedUpEntry(key);
- }
-
- public void clearLookedUpEntries() {
- xaAdapter.clearLookedUpEntries();
- }
-}
Copied: trunk/core/src/main/java/org/infinispan/context/impl/LocalTxInvocationContext.java (from rev 252, trunk/core/src/main/java/org/infinispan/context/impl/InitiatorTxInvocationContext.java)
===================================================================
--- trunk/core/src/main/java/org/infinispan/context/impl/LocalTxInvocationContext.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/context/impl/LocalTxInvocationContext.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -0,0 +1,74 @@
+package org.infinispan.context.impl;
+
+import org.infinispan.commands.write.WriteCommand;
+import org.infinispan.container.entries.CacheEntry;
+import org.infinispan.transaction.xa.GlobalTransaction;
+import org.infinispan.transaction.xa.TransactionXaAdapter;
+import org.infinispan.util.BidirectionalMap;
+
+import javax.transaction.Transaction;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * // TODO: Mircea: Document this!
+ *
+ * @author Mircea.Markus at jboss.com
+ * @since 4.0
+ */
+public class LocalTxInvocationContext extends AbstractTxInvocationContext {
+
+ private TransactionXaAdapter xaAdapter;
+
+ public Transaction getRunningTransaction() {
+ return xaAdapter.getTransaction();
+ }
+
+ public boolean isOriginLocal() {
+ return true;
+ }
+
+ public boolean isInTxScope() {
+ return true;
+ }
+
+ public Object getLockOwner() {
+ return xaAdapter.getGlobalTx();
+ }
+
+ public GlobalTransaction getGlobalTransaction() {
+ return xaAdapter.getGlobalTx();
+ }
+
+ public List<WriteCommand> getModifications() {
+ return xaAdapter.getModifications();
+ }
+
+ public void setXaCache(TransactionXaAdapter xaAdapter) {
+ this.xaAdapter = xaAdapter;
+ }
+
+ public CacheEntry lookupEntry(Object key) {
+ return xaAdapter != null ? xaAdapter.lookupEntry(key) : null;
+ }
+
+ public BidirectionalMap<Object, CacheEntry> getLookedUpEntries() {
+ return xaAdapter.getLookedUpEntries();
+ }
+
+ public void putLookedUpEntry(Object key, CacheEntry e) {
+ xaAdapter.putLookedUpEntry(key, e);
+ }
+
+ public void putLookedUpEntries(Map<Object, CacheEntry> lookedUpEntries) {
+ xaAdapter.putLookedUpEntries(lookedUpEntries);
+ }
+
+ public void removeLookedUpEntry(Object key) {
+ xaAdapter.removeLookedUpEntry(key);
+ }
+
+ public void clearLookedUpEntries() {
+ xaAdapter.clearLookedUpEntries();
+ }
+}
Property changes on: trunk/core/src/main/java/org/infinispan/context/impl/LocalTxInvocationContext.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Modified: trunk/core/src/main/java/org/infinispan/context/impl/RemoteTxInvocationContext.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/context/impl/RemoteTxInvocationContext.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/context/impl/RemoteTxInvocationContext.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -3,12 +3,10 @@
import org.infinispan.commands.write.WriteCommand;
import org.infinispan.container.entries.CacheEntry;
import org.infinispan.transaction.xa.GlobalTransaction;
-import org.infinispan.util.BidirectionalLinkedHashMap;
+import org.infinispan.transaction.xa.RemoteTransaction;
import org.infinispan.util.BidirectionalMap;
import javax.transaction.Transaction;
-import java.util.ArrayList;
-import java.util.Arrays;
import java.util.List;
import java.util.Map;
@@ -20,12 +18,9 @@
*/
public class RemoteTxInvocationContext extends AbstractTxInvocationContext {
- private List<WriteCommand> modifications;
- private BidirectionalMap<Object, CacheEntry> lookedUpEntries;
+ private RemoteTransaction remoteTransaction;
- protected GlobalTransaction tx;
-
public RemoteTxInvocationContext() {
}
@@ -34,11 +29,11 @@
}
public Object getLockOwner() {
- return tx;
+ return remoteTransaction.getGlobalTransaction();
}
- public GlobalTransaction getClusterTransactionId() {
- return tx;
+ public GlobalTransaction getGlobalTransaction() {
+ return remoteTransaction.getGlobalTransaction();
}
public boolean isInTxScope() {
@@ -50,74 +45,54 @@
}
public List<WriteCommand> getModifications() {
- return modifications;
+ return remoteTransaction.getModifications();
}
- public void initialize(WriteCommand[] modifications, GlobalTransaction tx) {
- this.modifications = Arrays.asList(modifications);
- lookedUpEntries = new BidirectionalLinkedHashMap<Object, CacheEntry>(modifications.length);
- this.tx = tx;
+ public void setRemoteTransaction(RemoteTransaction remoteTransaction) {
+ this.remoteTransaction = remoteTransaction;
}
public CacheEntry lookupEntry(Object key) {
- return lookedUpEntries.get(key);
+ return remoteTransaction.lookupEntry(key);
}
public BidirectionalMap<Object, CacheEntry> getLookedUpEntries() {
- return lookedUpEntries;
+ return remoteTransaction.getLookedUpEntries();
}
public void putLookedUpEntry(Object key, CacheEntry e) {
- lookedUpEntries.put(key, e);
+ remoteTransaction.putLookedUpEntry(key, e);
}
public void removeLookedUpEntry(Object key) {
- lookedUpEntries.remove(key);
+ remoteTransaction.removeLookedUpEntry(key);
}
public void clearLookedUpEntries() {
- lookedUpEntries.clear();
+ remoteTransaction.clearLookedUpEntries();
}
public void putLookedUpEntries(Map<Object, CacheEntry> lookedUpEntries) {
- lookedUpEntries.putAll(lookedUpEntries);
+ remoteTransaction.putLookedUpEntries(lookedUpEntries);
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (!(o instanceof RemoteTxInvocationContext)) return false;
-
- RemoteTxInvocationContext context = (RemoteTxInvocationContext) o;
- return tx.equals(context.tx);
+ RemoteTxInvocationContext that = (RemoteTxInvocationContext) o;
+ return remoteTransaction.equals(that.remoteTransaction);
}
@Override
public int hashCode() {
- return tx.hashCode();
+ return remoteTransaction.hashCode();
}
@Override
- public String toString() {
- return "RemoteTxInvocationContext{" +
- "modifications=" + modifications +
- ", lookedUpEntries=" + lookedUpEntries +
- ", tx=" + tx +
- "} " + super.toString();
- }
-
- @Override
public RemoteTxInvocationContext clone() {
RemoteTxInvocationContext dolly = (RemoteTxInvocationContext) super.clone();
- if (modifications != null) {
- dolly.modifications = new ArrayList<WriteCommand>(modifications);
- }
- if (lookedUpEntries != null) {
- dolly.lookedUpEntries = new BidirectionalLinkedHashMap<Object, CacheEntry>(lookedUpEntries);
- }
- if (tx != null) {
- dolly.tx = (GlobalTransaction) tx.clone();
- }
+ dolly.remoteTransaction = (RemoteTransaction) remoteTransaction.clone();
return dolly;
}
}
Modified: trunk/core/src/main/java/org/infinispan/context/impl/TxInvocationContext.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/context/impl/TxInvocationContext.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/context/impl/TxInvocationContext.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -21,7 +21,7 @@
Set<Address> getTransactionParticipants();
- GlobalTransaction getClusterTransactionId();
+ GlobalTransaction getGlobalTransaction();
List<WriteCommand> getModifications();
Modified: trunk/core/src/main/java/org/infinispan/factories/EmptyConstructorFactory.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/factories/EmptyConstructorFactory.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/factories/EmptyConstructorFactory.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -7,7 +7,7 @@
import org.infinispan.factories.scopes.Scopes;
import org.infinispan.notifications.cachemanagerlistener.CacheManagerNotifier;
import org.infinispan.remoting.InboundInvocationHandler;
-import org.infinispan.transaction.xa.TxEnlistingManager;
+import org.infinispan.transaction.xa.TransactionTable;
/**
* Factory for building global-scope components which have default empty constructors
@@ -16,7 +16,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, TxEnlistingManager.class})
+ at DefaultFactoryFor(classes = {InboundInvocationHandler.class, CacheManagerNotifier.class, RemoteCommandFactory.class, TransactionTable.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-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/factories/EmptyConstructorNamedCacheFactory.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -25,8 +25,8 @@
import org.infinispan.batch.BatchContainer;
import org.infinispan.commands.CommandsFactory;
import org.infinispan.config.ConfigurationException;
-import org.infinispan.context.container.InvocationContextContainer;
-import org.infinispan.context.container.ReplicationInvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainerImpl;
import org.infinispan.eviction.EvictionManager;
import org.infinispan.factories.annotations.DefaultFactoryFor;
import org.infinispan.loaders.CacheLoaderManager;
@@ -54,7 +54,7 @@
if (componentType.equals(Marshaller.class)) {
componentImpl = VersionAwareMarshaller.class;
} else if (componentType.equals(InvocationContextContainer.class)) {
- componentImpl = ReplicationInvocationContextContainer.class;
+ componentImpl = InvocationContextContainerImpl.class;
} else {
// add an "Impl" to the end of the class name and try again
componentImpl = getClass().getClassLoader().loadClass(componentType.getName() + "Impl");
Modified: trunk/core/src/main/java/org/infinispan/interceptors/BatchingInterceptor.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/interceptors/BatchingInterceptor.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/interceptors/BatchingInterceptor.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -24,7 +24,7 @@
import org.infinispan.batch.BatchContainer;
import org.infinispan.commands.VisitableCommand;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.factories.annotations.Inject;
import org.infinispan.interceptors.base.CommandInterceptor;
@@ -63,7 +63,7 @@
try {
transactionManager.resume(tx);
//this will make the call with a tx invocation context
- return invokeNextInterceptor(icc.getLocalInvocationContext(true), command);
+ return invokeNextInterceptor(icc.getLocalInvocationContext(), command);
} finally {
if (transactionManager.getTransaction() != null && batchContainer.isSuspendTxAfterInvocation())
transactionManager.suspend();
Modified: trunk/core/src/main/java/org/infinispan/interceptors/InterceptorChain.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/interceptors/InterceptorChain.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/interceptors/InterceptorChain.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -24,7 +24,7 @@
import org.infinispan.CacheException;
import org.infinispan.commands.VisitableCommand;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.factories.annotations.Inject;
import org.infinispan.factories.annotations.Start;
import org.infinispan.factories.scopes.Scope;
Modified: trunk/core/src/main/java/org/infinispan/interceptors/NotificationInterceptor.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/interceptors/NotificationInterceptor.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/interceptors/NotificationInterceptor.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -46,21 +46,21 @@
@Override
public Object visitPrepareCommand(TxInvocationContext ctx, PrepareCommand command) throws Throwable {
Object retval = invokeNextInterceptor(ctx, command);
- if (command.isOnePhaseCommit()) notifier.notifyTransactionCompleted(ctx.getClusterTransactionId(), true, ctx);
+ if (command.isOnePhaseCommit()) notifier.notifyTransactionCompleted(ctx.getGlobalTransaction(), true, ctx);
return retval;
}
@Override
public Object visitCommitCommand(TxInvocationContext ctx, CommitCommand command) throws Throwable {
Object retval = invokeNextInterceptor(ctx, command);
- notifier.notifyTransactionCompleted(ctx.getClusterTransactionId(), true, ctx);
+ notifier.notifyTransactionCompleted(ctx.getGlobalTransaction(), true, ctx);
return retval;
}
@Override
public Object visitRollbackCommand(TxInvocationContext ctx, RollbackCommand command) throws Throwable {
Object retval = invokeNextInterceptor(ctx, command);
- notifier.notifyTransactionCompleted(ctx.getClusterTransactionId(), false, ctx);
+ notifier.notifyTransactionCompleted(ctx.getGlobalTransaction(), false, ctx);
return retval;
}
}
Modified: trunk/core/src/main/java/org/infinispan/interceptors/TxInterceptor.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/interceptors/TxInterceptor.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/interceptors/TxInterceptor.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -17,16 +17,21 @@
import org.infinispan.commands.write.WriteCommand;
import org.infinispan.context.Flag;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.impl.InitiatorTxInvocationContext;
+import org.infinispan.context.impl.LocalTxInvocationContext;
import org.infinispan.context.impl.TxInvocationContext;
import org.infinispan.factories.annotations.Inject;
import org.infinispan.interceptors.base.CommandInterceptor;
import org.infinispan.jmx.annotations.ManagedAttribute;
import org.infinispan.jmx.annotations.ManagedOperation;
import org.infinispan.transaction.TransactionLog;
-import org.infinispan.transaction.xa.TxEnlistingManager;
import org.infinispan.transaction.xa.TransactionXaAdapter;
+import org.infinispan.transaction.xa.TransactionTable;
+import javax.transaction.RollbackException;
+import javax.transaction.Status;
+import javax.transaction.SystemException;
+import javax.transaction.Transaction;
+import javax.transaction.TransactionManager;
import java.util.Arrays;
import java.util.concurrent.atomic.AtomicLong;
@@ -39,8 +44,9 @@
*/
public class TxInterceptor extends CommandInterceptor {
- private TxEnlistingManager txEnlistingManager;
+ private TransactionManager tm;
private TransactionLog transactionLog;
+ private TransactionTable txTable;
private final AtomicLong prepares = new AtomicLong(0);
private final AtomicLong commits = new AtomicLong(0);
@@ -49,9 +55,10 @@
@Inject
- public void init(TxEnlistingManager txEnlistingManager, TransactionLog transactionLog) {
- this.txEnlistingManager = txEnlistingManager;
+ public void init(TransactionManager tm, TransactionTable txTable, TransactionLog transactionLog) {
+ this.tm = tm;
this.transactionLog = transactionLog;
+ this.txTable = txTable;
setStatisticsEnabled(configuration.isExposeJmxStatistics());
}
@@ -69,7 +76,7 @@
if (!command.isOnePhaseCommit()) {
transactionLog.logPrepare(command);
} else {
- transactionLog.logOnePhaseCommit(ctx.getClusterTransactionId(), Arrays.asList(command.getModifications()));
+ transactionLog.logOnePhaseCommit(ctx.getGlobalTransaction(), Arrays.asList(command.getModifications()));
}
if (getStatisticsEnabled()) prepares.incrementAndGet();
return invokeNextInterceptor(ctx, command);
@@ -88,12 +95,12 @@
transactionLog.rollback(command.getGlobalTransaction());
return invokeNextInterceptor(ctx, command);
}
-
+
@Override
- public Object visitLockControlCommand(InvocationContext ctx, LockControlCommand command) throws Throwable{
+ public Object visitLockControlCommand(InvocationContext ctx, LockControlCommand command) throws Throwable {
return enlistReadAndInvokeNext(ctx, command);
}
-
+
/**
* Designed to be overridden. Returns a VisitableCommand fit for replaying locally, based on the modification passed
* in. If a null value is returned, this means that the command should not be replayed.
@@ -153,27 +160,40 @@
private Object enlistReadAndInvokeNext(InvocationContext ctx, VisitableCommand command) throws Throwable {
if (shouldEnlist(ctx)) {
- TransactionXaAdapter xaAdapter = txEnlistingManager.enlist(ctx);
- InitiatorTxInvocationContext initiatorTxContext = (InitiatorTxInvocationContext) ctx;
- initiatorTxContext.setXaCache(xaAdapter);
+ TransactionXaAdapter xaAdapter = enlist();
+ LocalTxInvocationContext localTxContext = (LocalTxInvocationContext) ctx;
+ localTxContext.setXaCache(xaAdapter);
}
return invokeNextInterceptor(ctx, command);
}
private Object enlistWriteAndInvokeNext(InvocationContext ctx, WriteCommand command) throws Throwable {
if (shouldEnlist(ctx)) {
- TransactionXaAdapter xaAdapter = txEnlistingManager.enlist(ctx);
- InitiatorTxInvocationContext initiatorTxContext = (InitiatorTxInvocationContext) ctx;
+ TransactionXaAdapter xaAdapter = enlist();
+ LocalTxInvocationContext localTxContext = (LocalTxInvocationContext) ctx;
if (!isLocalModeForced(ctx)) {
xaAdapter.addModification(command);
}
- initiatorTxContext.setXaCache(xaAdapter);
+ localTxContext.setXaCache(xaAdapter);
}
if (!ctx.isInTxScope())
transactionLog.logNoTxWrite(command);
return invokeNextInterceptor(ctx, command);
}
+ public TransactionXaAdapter enlist() throws SystemException, RollbackException {
+ Transaction transaction = tm.getTransaction();
+ if (transaction == null) throw new IllegalStateException("This should only be called in an tx scope");
+ int status = transaction.getStatus();
+ if (!isValid(status)) throw new IllegalStateException("Transaction " + transaction +
+ " is not in a valid state to be invoking cache operations on.");
+ return txTable.getOrCreateXaAdapter(transaction);
+ }
+
+ private boolean isValid(int status) {
+ return status == Status.STATUS_ACTIVE || status == Status.STATUS_PREPARING;
+ }
+
private boolean shouldEnlist(InvocationContext ctx) {
return ctx.isInTxScope() & ctx.isOriginLocal();
}
Modified: trunk/core/src/main/java/org/infinispan/notifications/cachelistener/CacheNotifierImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/notifications/cachelistener/CacheNotifierImpl.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/notifications/cachelistener/CacheNotifierImpl.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -23,7 +23,7 @@
import org.infinispan.Cache;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.context.impl.TxInvocationContext;
import org.infinispan.factories.annotations.Inject;
import org.infinispan.notifications.AbstractListenerImpl;
@@ -241,7 +241,7 @@
private void setTx(InvocationContext ctx, EventImpl e) {
if (ctx.isInTxScope()) {
- GlobalTransaction tx = ((TxInvocationContext) ctx).getClusterTransactionId();
+ GlobalTransaction tx = ((TxInvocationContext) ctx).getGlobalTransaction();
e.setTransactionId(tx);
}
}
Modified: trunk/core/src/main/java/org/infinispan/remoting/InboundInvocationHandlerImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/remoting/InboundInvocationHandlerImpl.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/remoting/InboundInvocationHandlerImpl.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -3,7 +3,7 @@
import org.infinispan.commands.CommandsFactory;
import org.infinispan.commands.remote.CacheRpcCommand;
import org.infinispan.config.Configuration;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.factories.ComponentRegistry;
import org.infinispan.factories.GlobalComponentRegistry;
import org.infinispan.factories.annotations.Inject;
Modified: trunk/core/src/main/java/org/infinispan/statetransfer/StateTransferManagerImpl.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/statetransfer/StateTransferManagerImpl.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/statetransfer/StateTransferManagerImpl.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -30,8 +30,8 @@
import org.infinispan.container.DataContainer;
import org.infinispan.container.entries.InternalCacheEntry;
import org.infinispan.context.Flag;
-import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
+import org.infinispan.context.impl.RemoteTxInvocationContext;
import org.infinispan.factories.annotations.Inject;
import org.infinispan.factories.annotations.Start;
import org.infinispan.interceptors.InterceptorChain;
@@ -46,6 +46,8 @@
import org.infinispan.remoting.transport.Address;
import org.infinispan.remoting.transport.DistributedSync;
import org.infinispan.transaction.TransactionLog;
+import org.infinispan.transaction.xa.RemoteTransaction;
+import org.infinispan.transaction.xa.TransactionTable;
import org.infinispan.util.Util;
import org.infinispan.util.logging.Log;
import org.infinispan.util.logging.LogFactory;
@@ -74,6 +76,7 @@
InvocationContextContainer invocationContextContainer;
InterceptorChain interceptorChain;
CommandsFactory commandsFactory;
+ TransactionTable txTable;
private static final Log log = LogFactory.getLog(StateTransferManagerImpl.class);
private static final boolean trace = log.isTraceEnabled();
private static final Byte DELIMITER = (byte) 123;
@@ -86,7 +89,7 @@
public void injectDependencies(RpcManager rpcManager, AdvancedCache cache, Configuration configuration,
DataContainer dataContainer, CacheLoaderManager clm, Marshaller marshaller,
TransactionLog transactionLog, InterceptorChain interceptorChain, InvocationContextContainer invocationContextContainer,
- CommandsFactory commandsFactory) {
+ CommandsFactory commandsFactory, TransactionTable txTable) {
this.rpcManager = rpcManager;
this.cache = cache;
this.configuration = configuration;
@@ -97,6 +100,7 @@
this.invocationContextContainer = invocationContextContainer;
this.interceptorChain = interceptorChain;
this.commandsFactory = commandsFactory;
+ this.txTable = txTable;
}
@Start(priority = 55)
@@ -214,13 +218,15 @@
private void processCommitLog(ObjectInput oi) throws Exception {
if (trace) log.trace("Applying commit log");
Object object = marshaller.objectFromObjectStream(oi);
+ RemoteTxInvocationContext ctx = invocationContextContainer.getRemoteTxInvocationContext();
while (object instanceof TransactionLog.LogEntry) {
TransactionLog.LogEntry logEntry = (TransactionLog.LogEntry) object;
+ RemoteTransaction remoteTransaction = txTable.getRemoteTransaction(logEntry.getTransaction());
+ ctx.setRemoteTransaction(remoteTransaction);
WriteCommand[] mods = logEntry.getModifications();
if (trace) log.trace("Mods = {0}", Arrays.toString(mods));
for (WriteCommand mod : mods) {
commandsFactory.initializeReplicableCommand(mod);
- InvocationContext ctx = invocationContextContainer.getRemoteTxInvocationContext(logEntry.getTransaction(),true);
ctx.setFlags(Flag.CACHE_MODE_LOCAL, Flag.SKIP_CACHE_STATUS_CHECK);
interceptorChain.invoke(ctx, mod);
}
@@ -254,7 +260,9 @@
if (!transactionLog.hasPendingPrepare(command)) {
if (trace) log.trace("Applying pending prepare {0}", command);
commandsFactory.initializeReplicableCommand(command);
- InvocationContext ctx = invocationContextContainer.getRemoteTxInvocationContext(command.getGlobalTransaction(), true);
+ RemoteTxInvocationContext ctx = invocationContextContainer.getRemoteTxInvocationContext();
+ RemoteTransaction transaction = txTable.createRemoteTransaction(command.getGlobalTransaction(), command.getModifications());
+ ctx.setRemoteTransaction(transaction);
ctx.setFlags(Flag.CACHE_MODE_LOCAL, Flag.SKIP_CACHE_STATUS_CHECK);
interceptorChain.invoke(ctx, command);
} else {
Added: trunk/core/src/main/java/org/infinispan/transaction/xa/CacheTransaction.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/transaction/xa/CacheTransaction.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/transaction/xa/CacheTransaction.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -0,0 +1,32 @@
+package org.infinispan.transaction.xa;
+
+import org.infinispan.commands.write.WriteCommand;
+import org.infinispan.container.entries.CacheEntry;
+import org.infinispan.util.BidirectionalMap;
+
+import java.util.List;
+import java.util.Map;
+
+/**
+ * // TODO: Mircea: Document this!
+ *
+ * @author
+ */
+public interface CacheTransaction {
+
+ public GlobalTransaction getGlobalTransaction();
+
+ public List<WriteCommand> getModifications();
+
+ public CacheEntry lookupEntry(Object key);
+
+ public BidirectionalMap<Object, CacheEntry> getLookedUpEntries();
+
+ public void putLookedUpEntry(Object key, CacheEntry e);
+
+ public void putLookedUpEntries(Map<Object, CacheEntry> lookedUpEntries);
+
+ public void removeLookedUpEntry(Object key);
+
+ public void clearLookedUpEntries();
+}
Property changes on: trunk/core/src/main/java/org/infinispan/transaction/xa/CacheTransaction.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Added: trunk/core/src/main/java/org/infinispan/transaction/xa/RemoteTransaction.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/transaction/xa/RemoteTransaction.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/transaction/xa/RemoteTransaction.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -0,0 +1,98 @@
+package org.infinispan.transaction.xa;
+
+import org.infinispan.commands.write.WriteCommand;
+import org.infinispan.container.entries.CacheEntry;
+import org.infinispan.util.BidirectionalLinkedHashMap;
+import org.infinispan.util.BidirectionalMap;
+
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * // TODO: Mircea: Document this!
+ *
+ * @author
+ */
+public class RemoteTransaction implements CacheTransaction, Cloneable {
+
+ private List<WriteCommand> modifications;
+
+ private BidirectionalLinkedHashMap<Object, CacheEntry> lookedUpEntries;
+
+ private GlobalTransaction tx;
+
+
+ public RemoteTransaction(WriteCommand[] modifications, GlobalTransaction tx) {
+ this.modifications = Arrays.asList(modifications);
+ lookedUpEntries = new BidirectionalLinkedHashMap<Object, CacheEntry>(modifications.length);
+ this.tx = tx;
+ }
+
+ public GlobalTransaction getGlobalTransaction() {
+ return tx;
+ }
+
+ public List<WriteCommand> getModifications() {
+ return modifications;
+ }
+
+ public CacheEntry lookupEntry(Object key) {
+ return lookedUpEntries.get(key);
+ }
+
+ public BidirectionalMap<Object, CacheEntry> getLookedUpEntries() {
+ return lookedUpEntries;
+ }
+
+ public void putLookedUpEntry(Object key, CacheEntry e) {
+ lookedUpEntries.put(key, e);
+ }
+
+ public void putLookedUpEntries(Map<Object, CacheEntry> lookedUpEntries) {
+ lookedUpEntries.putAll(lookedUpEntries);
+ }
+
+ public void removeLookedUpEntry(Object key) {
+ lookedUpEntries.remove(key);
+ }
+
+ public void clearLookedUpEntries() {
+ lookedUpEntries.clear();
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) return true;
+ if (!(o instanceof RemoteTransaction)) return false;
+ RemoteTransaction that = (RemoteTransaction) o;
+ return tx.equals(that.tx);
+ }
+
+ @Override
+ public int hashCode() {
+ return tx.hashCode();
+ }
+
+ @Override
+ public Object clone() {
+ try {
+ RemoteTransaction dolly = (RemoteTransaction) super.clone();
+ dolly.modifications = new ArrayList<WriteCommand>(modifications);
+ dolly.lookedUpEntries = lookedUpEntries.clone();
+ return dolly;
+ } catch (CloneNotSupportedException e) {
+ throw new IllegalStateException("Impossible!!");
+ }
+ }
+
+ @Override
+ public String toString() {
+ return "RemoteTransaction{" +
+ "modifications=" + modifications +
+ ", lookedUpEntries=" + lookedUpEntries +
+ ", tx=" + tx +
+ '}';
+ }
+}
Property changes on: trunk/core/src/main/java/org/infinispan/transaction/xa/RemoteTransaction.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Added: trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionTable.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionTable.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionTable.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -0,0 +1,108 @@
+package org.infinispan.transaction.xa;
+
+import org.infinispan.CacheException;
+import org.infinispan.commands.CommandsFactory;
+import org.infinispan.commands.write.WriteCommand;
+import org.infinispan.config.Configuration;
+import org.infinispan.context.InvocationContextContainer;
+import org.infinispan.factories.annotations.Inject;
+import org.infinispan.interceptors.InterceptorChain;
+import org.infinispan.notifications.cachelistener.CacheNotifier;
+import org.infinispan.remoting.rpc.CacheRpcManager;
+import org.infinispan.remoting.transport.Address;
+import org.infinispan.util.logging.Log;
+import org.infinispan.util.logging.LogFactory;
+
+import javax.transaction.Transaction;
+import java.util.HashMap;
+import java.util.Map;
+
+/**
+ * // TODO: Mircea: Document this!
+ *
+ * @author
+ */
+public class TransactionTable {
+
+ private static Log log = LogFactory.getLog(TransactionTable.class);
+
+ private Map<Transaction, TransactionXaAdapter> localTransactions = new HashMap<Transaction, TransactionXaAdapter>();
+
+ private Map<GlobalTransaction, RemoteTransaction> remoteTransactions = new HashMap<GlobalTransaction, RemoteTransaction>();
+
+ private CommandsFactory commandsFactory;
+ private Configuration configuration;
+ private InvocationContextContainer icc;
+ private InterceptorChain invoker;
+ private CacheNotifier notifier;
+ private CacheRpcManager rpcManager;
+
+
+ @Inject
+ public void initialize(CommandsFactory commandsFactory, CacheRpcManager rpcManager, Configuration configuration,
+ InvocationContextContainer icc, InterceptorChain invoker, CacheNotifier notifier) {
+ this.commandsFactory = commandsFactory;
+ this.rpcManager = rpcManager;
+ this.configuration = configuration;
+ this.icc = icc;
+ this.invoker = invoker;
+ this.notifier = notifier;
+ }
+
+
+ public RemoteTransaction getRemoteTransaction(GlobalTransaction txId) {
+ return remoteTransactions.get(txId);
+ }
+
+ public RemoteTransaction createRemoteTransaction(GlobalTransaction globalTx, WriteCommand[] modifications) {
+ RemoteTransaction remoteTransaction = new RemoteTransaction(modifications, globalTx);
+ RemoteTransaction transaction = remoteTransactions.put(globalTx, remoteTransaction);
+ if (transaction != null) {
+ String message = "A remote transaction with the given id was already registred!!!";
+ log.error(message);
+ throw new IllegalStateException(message);
+ }
+ if (log.isTraceEnabled()) {
+ log.trace("Created and regostered tremote transaction " + remoteTransaction);
+ }
+ return remoteTransaction;
+ }
+
+ public TransactionXaAdapter getOrCreateXaAdapter(Transaction transaction) {
+ TransactionXaAdapter current = localTransactions.get(transaction);
+ if (current == null) {
+ Address localAddress = rpcManager != null ? rpcManager.getLocalAddress() : null;
+ GlobalTransaction tx = localAddress == null ? new GlobalTransaction(false) : new GlobalTransaction(localAddress, false);
+ current = new TransactionXaAdapter(tx, icc, invoker, commandsFactory, configuration, this, transaction);
+ localTransactions.put(transaction, current);
+ try {
+ transaction.enlistResource(current);
+ } catch (Exception e) {
+ log.error("Failed to emlist TransactionXaAdapter to transaction");
+ throw new CacheException(e);
+ }
+ notifier.notifyTransactionRegistered(tx, icc.getThreadContext());
+ }
+ return current;
+ }
+
+ public boolean removeLocalTransaction(Transaction tx) {
+ return localTransactions.remove(tx) != null;
+ }
+
+ public boolean removeRemoteTransaction(GlobalTransaction txId) {
+ return remoteTransactions.remove(txId) != null;
+ }
+
+ public int getRemoteTxCount() {
+ return remoteTransactions.size();
+ }
+
+ public int getLocalTxCount() {
+ return localTransactions.size();
+ }
+
+ public TransactionXaAdapter getXaCacheAdapter(Transaction tx) {
+ return localTransactions.get(tx);
+ }
+}
Property changes on: trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionTable.java
___________________________________________________________________
Name: svn:keywords
+ Id Revision
Name: svn:eol-style
+ LF
Deleted: trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionXaAdapter.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionXaAdapter.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionXaAdapter.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -1,214 +0,0 @@
-package org.infinispan.transaction.xa;
-
-import org.infinispan.commands.CommandsFactory;
-import org.infinispan.commands.tx.CommitCommand;
-import org.infinispan.commands.tx.PrepareCommand;
-import org.infinispan.commands.tx.RollbackCommand;
-import org.infinispan.commands.write.WriteCommand;
-import org.infinispan.config.Configuration;
-import org.infinispan.container.entries.CacheEntry;
-import org.infinispan.context.container.InvocationContextContainer;
-import org.infinispan.context.impl.InitiatorTxInvocationContext;
-import org.infinispan.interceptors.InterceptorChain;
-import org.infinispan.util.BidirectionalLinkedHashMap;
-import org.infinispan.util.BidirectionalMap;
-import org.infinispan.util.InfinispanCollections;
-import org.infinispan.util.logging.Log;
-import org.infinispan.util.logging.LogFactory;
-
-import javax.transaction.Transaction;
-import javax.transaction.xa.XAException;
-import javax.transaction.xa.XAResource;
-import javax.transaction.xa.Xid;
-import java.util.ArrayList;
-import java.util.List;
-import java.util.Map;
-
-/**
- * // TODO: Mircea: Document this!
- *
- * @author Mircea.Markus at jboss.com
- * @since 4.0
- */
-public class TransactionXaAdapter implements XAResource {
-
- private static Log log = LogFactory.getLog(TransactionXaAdapter.class);
-
- private int txTimeout;
-
- private List<WriteCommand> modifications;
- private BidirectionalMap<Object, CacheEntry> lookedUpEntries;
-
- private GlobalTransaction transactionIdentifier;
- private InvocationContextContainer icc;
- private InterceptorChain invoker;
-
- private CommandsFactory commandsFactory;
- private Configuration configuration;
-
- private TxEnlistingManager txEnlistingManager;
- private Transaction transaction;
-
- public TransactionXaAdapter(GlobalTransaction transactionIdentifier, InvocationContextContainer icc, InterceptorChain invoker,
- CommandsFactory commandsFactory, Configuration configuration, TxEnlistingManager txEnlistingManager,
- Transaction transaction) {
- this.transactionIdentifier = transactionIdentifier;
- this.icc = icc;
- this.invoker = invoker;
- this.commandsFactory = commandsFactory;
- this.configuration = configuration;
- this.txEnlistingManager = txEnlistingManager;
- this.transaction = transaction;
- }
-
- public void addModification(WriteCommand mod) {
- if (modifications == null) {
- modifications = new ArrayList<WriteCommand>(8);
- }
- modifications.add(mod);
- }
-
- public int prepare(Xid xid) throws XAException {
- if (configuration.isOnePhaseCommit()) {
- if (log.isTraceEnabled())
- log.trace("Recieved prepare for tx: " + xid + " . Skipping call as 1PC will be used.");
- return XA_OK;
- }
-
- PrepareCommand prepareCommand = commandsFactory.buildPrepareCommand(transactionIdentifier, modifications, configuration.isOnePhaseCommit());
- if (log.isTraceEnabled()) log.trace("Sending prepare command through the chain: " + prepareCommand);
-
- InitiatorTxInvocationContext ctx = icc.getInitiatorTxInvocationContext();
- ctx.setXaCache(this);
- try {
- invoker.invoke(ctx, prepareCommand);
- return XA_OK; //todo validate code here
- } catch (Throwable e) {
- log.error("Error while processing PrepareCommand", e);
- throw new XAException(XAException.XAER_RMERR);//todo validate code here
- }
- }
-
- public void commit(Xid xid, boolean b) throws XAException {
- if (log.isTraceEnabled()) log.trace("commiting TransactionXaAdapter: " + transactionIdentifier);
- try {
- InitiatorTxInvocationContext ctx = icc.getInitiatorTxInvocationContext();
- ctx.setXaCache(this);
- if (configuration.isOnePhaseCommit()) {
- if (log.isTraceEnabled()) log.trace("Doing an 1PC prepare call on the interceptor chain");
- PrepareCommand command = commandsFactory.buildPrepareCommand(transactionIdentifier, modifications, true);
- try {
- invoker.invoke(ctx, command);
- } catch (Throwable e) {
- log.error("Error while processing 1PC PrepareCommand", e);
- throw new XAException(XAException.XAER_RMERR);
- }
- } else {
- CommitCommand commitCommand = commandsFactory.buildCommitCommand(transactionIdentifier);
- try {
- invoker.invoke(ctx, commitCommand);
- } catch (Throwable e) {
- log.error("Error while processing 1PC PrepareCommand", e);
- throw new XAException(XAException.XAER_RMERR);
- }
- }
- } finally {
- txEnlistingManager.delist(transaction);
- this.modifications = null;
- }
- }
-
- public void rollback(Xid xid) throws XAException {
- RollbackCommand rollbackCommand = commandsFactory.buildRollbackCommand(transactionIdentifier);
- InitiatorTxInvocationContext ctx = icc.getInitiatorTxInvocationContext();
- ctx.setXaCache(this);
- try {
- invoker.invoke(ctx, rollbackCommand);
- } catch (Throwable e) {
- log.error("Exception while ", e);
- throw new XAException(XAException.XA_HEURHAZ);
- } finally {
- txEnlistingManager.delist(transaction);
- this.modifications = null;
- }
- }
-
- public void start(Xid xid, int i) throws XAException {
- if (log.isTraceEnabled()) log.trace("start called");
- }
-
- public void end(Xid xid, int i) throws XAException {
- if (log.isTraceEnabled()) log.trace("end called");
- }
-
- public void forget(Xid xid) throws XAException {
- if (log.isTraceEnabled()) log.trace("forget called");
- }
-
- public int getTransactionTimeout() throws XAException {
- if (log.isTraceEnabled()) log.trace("start called");
- return txTimeout;
- }
-
- public boolean isSameRM(XAResource xaResource) throws XAException {
- if (!(xaResource instanceof TransactionXaAdapter)) {
- return false;
- }
- TransactionXaAdapter other = (TransactionXaAdapter) xaResource;
- return other.transactionIdentifier.equals(this.transactionIdentifier);
- }
-
- public Xid[] recover(int i) throws XAException {
- if (log.isTraceEnabled()) log.trace("recover called: " + i);
- return null; //todo validate with javadoc
- }
-
- public boolean setTransactionTimeout(int i) throws XAException {
- this.txTimeout = i;
- return true; //todo check javadoc
- }
-
- public void putLookedUpEntries(Map<Object, CacheEntry> entries) {
- initLookedUpEntries();
- lookedUpEntries.putAll(entries);
- }
-
- public CacheEntry lookupEntry(Object key) {
- if (lookedUpEntries == null) return null;
- return lookedUpEntries.get(key);
- }
-
- public BidirectionalMap<Object, CacheEntry> getLookedUpEntries() {
- return (BidirectionalMap<Object, CacheEntry>)
- (lookedUpEntries == null ? InfinispanCollections.emptyBidirectionalMap() : lookedUpEntries);
- }
-
- public void putLookedUpEntry(Object key, CacheEntry e) {
- initLookedUpEntries();
- lookedUpEntries.put(key, e);
- }
-
- public void removeLookedUpEntry(Object key) {
- if (lookedUpEntries != null) lookedUpEntries.remove(key);
- }
-
- public void clearLookedUpEntries() {
- if (lookedUpEntries != null) lookedUpEntries.clear();
- }
-
- private void initLookedUpEntries() {
- if (lookedUpEntries == null) lookedUpEntries = new BidirectionalLinkedHashMap<Object, CacheEntry>(4);
- }
-
- public GlobalTransaction getTransactionIdentifier() {
- return transactionIdentifier;
- }
-
- public List<WriteCommand> getModifications() {
- return modifications;
- }
-
- public Transaction getTransaction() {
- return transaction;
- }
-}
Added: trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionXaAdapter.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionXaAdapter.java (rev 0)
+++ trunk/core/src/main/java/org/infinispan/transaction/xa/TransactionXaAdapter.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -0,0 +1,218 @@
+package org.infinispan.transaction.xa;
+
+import org.infinispan.commands.CommandsFactory;
+import org.infinispan.commands.tx.CommitCommand;
+import org.infinispan.commands.tx.PrepareCommand;
+import org.infinispan.commands.tx.RollbackCommand;
+import org.infinispan.commands.write.WriteCommand;
+import org.infinispan.config.Configuration;
+import org.infinispan.container.entries.CacheEntry;
+import org.infinispan.context.InvocationContextContainer;
+import org.infinispan.context.impl.LocalTxInvocationContext;
+import org.infinispan.interceptors.InterceptorChain;
+import org.infinispan.util.BidirectionalLinkedHashMap;
+import org.infinispan.util.BidirectionalMap;
+import org.infinispan.util.InfinispanCollections;
+import org.infinispan.util.logging.Log;
+import org.infinispan.util.logging.LogFactory;
+
+import javax.transaction.Transaction;
+import javax.transaction.xa.XAException;
+import javax.transaction.xa.XAResource;
+import javax.transaction.xa.Xid;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+
+/**
+ * // TODO: Mircea: Document this!
+ *
+ * @author Mircea.Markus at jboss.com
+ * @since 4.0
+ */
+public class TransactionXaAdapter implements CacheTransaction, XAResource {
+
+ private static Log log = LogFactory.getLog(TransactionXaAdapter.class);
+
+ private int txTimeout;
+
+ private List<WriteCommand> modifications;
+ private BidirectionalMap<Object, CacheEntry> lookedUpEntries;
+
+ private GlobalTransaction globalTx;
+ private InvocationContextContainer icc;
+ private InterceptorChain invoker;
+
+ private CommandsFactory commandsFactory;
+ private Configuration configuration;
+
+ private TransactionTable txTable;
+ private Transaction transaction;
+
+ public TransactionXaAdapter(GlobalTransaction globalTx, InvocationContextContainer icc, InterceptorChain invoker,
+ CommandsFactory commandsFactory, Configuration configuration, TransactionTable txTable,
+ Transaction transaction) {
+ this.globalTx = globalTx;
+ this.icc = icc;
+ this.invoker = invoker;
+ this.commandsFactory = commandsFactory;
+ this.configuration = configuration;
+ this.txTable = txTable;
+ this.transaction = transaction;
+ }
+
+ public void addModification(WriteCommand mod) {
+ if (modifications == null) {
+ modifications = new ArrayList<WriteCommand>(8);
+ }
+ modifications.add(mod);
+ }
+
+ public int prepare(Xid xid) throws XAException {
+ if (configuration.isOnePhaseCommit()) {
+ if (log.isTraceEnabled())
+ log.trace("Recieved prepare for tx: " + xid + " . Skipping call as 1PC will be used.");
+ return XA_OK;
+ }
+
+ PrepareCommand prepareCommand = commandsFactory.buildPrepareCommand(globalTx, modifications, configuration.isOnePhaseCommit());
+ if (log.isTraceEnabled()) log.trace("Sending prepare command through the chain: " + prepareCommand);
+
+ LocalTxInvocationContext ctx = icc.getInitiatorTxInvocationContext();
+ ctx.setXaCache(this);
+ try {
+ invoker.invoke(ctx, prepareCommand);
+ return XA_OK; //todo validate code here
+ } catch (Throwable e) {
+ log.error("Error while processing PrepareCommand", e);
+ throw new XAException(XAException.XAER_RMERR);//todo validate code here
+ }
+ }
+
+ public void commit(Xid xid, boolean b) throws XAException {
+ if (log.isTraceEnabled()) log.trace("commiting TransactionXaAdapter: " + globalTx);
+ try {
+ LocalTxInvocationContext ctx = icc.getInitiatorTxInvocationContext();
+ ctx.setXaCache(this);
+ if (configuration.isOnePhaseCommit()) {
+ if (log.isTraceEnabled()) log.trace("Doing an 1PC prepare call on the interceptor chain");
+ PrepareCommand command = commandsFactory.buildPrepareCommand(globalTx, modifications, true);
+ try {
+ invoker.invoke(ctx, command);
+ } catch (Throwable e) {
+ log.error("Error while processing 1PC PrepareCommand", e);
+ throw new XAException(XAException.XAER_RMERR);
+ }
+ } else {
+ CommitCommand commitCommand = commandsFactory.buildCommitCommand(globalTx);
+ try {
+ invoker.invoke(ctx, commitCommand);
+ } catch (Throwable e) {
+ log.error("Error while processing 1PC PrepareCommand", e);
+ throw new XAException(XAException.XAER_RMERR);
+ }
+ }
+ } finally {
+ txTable.removeLocalTransaction(transaction);
+ this.modifications = null;
+ }
+ }
+
+ public void rollback(Xid xid) throws XAException {
+ RollbackCommand rollbackCommand = commandsFactory.buildRollbackCommand(globalTx);
+ LocalTxInvocationContext ctx = icc.getInitiatorTxInvocationContext();
+ ctx.setXaCache(this);
+ try {
+ invoker.invoke(ctx, rollbackCommand);
+ } catch (Throwable e) {
+ log.error("Exception while ", e);
+ throw new XAException(XAException.XA_HEURHAZ);
+ } finally {
+ txTable.removeLocalTransaction(transaction);
+ this.modifications = null;
+ }
+ }
+
+ public void start(Xid xid, int i) throws XAException {
+ if (log.isTraceEnabled()) log.trace("start called");
+ }
+
+ public void end(Xid xid, int i) throws XAException {
+ if (log.isTraceEnabled()) log.trace("end called");
+ }
+
+ public void forget(Xid xid) throws XAException {
+ if (log.isTraceEnabled()) log.trace("forget called");
+ }
+
+ public int getTransactionTimeout() throws XAException {
+ if (log.isTraceEnabled()) log.trace("start called");
+ return txTimeout;
+ }
+
+ public boolean isSameRM(XAResource xaResource) throws XAException {
+ if (!(xaResource instanceof TransactionXaAdapter)) {
+ return false;
+ }
+ TransactionXaAdapter other = (TransactionXaAdapter) xaResource;
+ return other.globalTx.equals(this.globalTx);
+ }
+
+ public Xid[] recover(int i) throws XAException {
+ if (log.isTraceEnabled()) log.trace("recover called: " + i);
+ return null; //todo validate with javadoc
+ }
+
+ public boolean setTransactionTimeout(int i) throws XAException {
+ this.txTimeout = i;
+ return true; //todo check javadoc
+ }
+
+ public void putLookedUpEntries(Map<Object, CacheEntry> entries) {
+ initLookedUpEntries();
+ lookedUpEntries.putAll(entries);
+ }
+
+ public CacheEntry lookupEntry(Object key) {
+ if (lookedUpEntries == null) return null;
+ return lookedUpEntries.get(key);
+ }
+
+ public BidirectionalMap<Object, CacheEntry> getLookedUpEntries() {
+ return (BidirectionalMap<Object, CacheEntry>)
+ (lookedUpEntries == null ? InfinispanCollections.emptyBidirectionalMap() : lookedUpEntries);
+ }
+
+ public void putLookedUpEntry(Object key, CacheEntry e) {
+ initLookedUpEntries();
+ lookedUpEntries.put(key, e);
+ }
+
+ private void initLookedUpEntries() {
+ if (lookedUpEntries == null) lookedUpEntries = new BidirectionalLinkedHashMap<Object, CacheEntry>(4);
+ }
+
+ public GlobalTransaction getGlobalTx() {
+ return globalTx;
+ }
+
+ public List<WriteCommand> getModifications() {
+ return modifications;
+ }
+
+ public Transaction getTransaction() {
+ return transaction;
+ }
+
+ public GlobalTransaction getGlobalTransaction() {
+ return globalTx;
+ }
+
+ public void removeLookedUpEntry(Object key) {
+ if (lookedUpEntries != null) lookedUpEntries.remove(key);
+ }
+
+ public void clearLookedUpEntries() {
+ if (lookedUpEntries != null) lookedUpEntries.clear();
+ }
+}
Deleted: trunk/core/src/main/java/org/infinispan/transaction/xa/TxEnlistingManager.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/transaction/xa/TxEnlistingManager.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/transaction/xa/TxEnlistingManager.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -1,154 +0,0 @@
-package org.infinispan.transaction.xa;
-
-import org.infinispan.CacheException;
-import org.infinispan.commands.CommandsFactory;
-import org.infinispan.config.Configuration;
-import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
-import org.infinispan.context.impl.RemoteTxInvocationContext;
-import org.infinispan.factories.annotations.Inject;
-import org.infinispan.interceptors.InterceptorChain;
-import org.infinispan.notifications.cachelistener.CacheNotifier;
-import org.infinispan.remoting.rpc.CacheRpcManager;
-
-import javax.transaction.RollbackException;
-import javax.transaction.Status;
-import javax.transaction.SystemException;
-import javax.transaction.Transaction;
-import javax.transaction.TransactionManager;
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
-
-/**
- * // TODO: Mircea: Document this!
- *
- * @author Mircea.Markus at jboss.com
- * @since 4.0
- */
-public class TxEnlistingManager {
- private TransactionManager tm;
- private InvocationContextContainer icc;
- private InterceptorChain invoker;
- private CacheNotifier notifier;
-
- private Map<Transaction, TransactionXaAdapter> tx2XaCacheMapping = new ConcurrentHashMap<Transaction, TransactionXaAdapter>();
- private Map<GlobalTransaction, RemoteTxInvocationContext> remoteTxMap;
- private CommandsFactory commandsFactory;
- private CacheRpcManager rpcManager;
- private Configuration configuration;
-
- @Inject
- public void initialize(CommandsFactory commandsFactory, CacheRpcManager rpcManager, Configuration configuration,
- InvocationContextContainer icc, InterceptorChain invoker,
- TransactionManager tm, CacheNotifier notifier) {
- this.commandsFactory = commandsFactory;
- this.rpcManager = rpcManager;
- this.configuration = configuration;
- this.tm = tm;
- this.icc = icc;
- this.invoker = invoker;
- this.notifier = notifier;
- }
-
- public TransactionXaAdapter enlist(InvocationContext ctx) throws SystemException, RollbackException {
- Transaction transaction = tm.getTransaction();
- if (transaction == null) throw new IllegalStateException("This should only be called in an tx scope");
- if (!isValid(transaction)) throw new IllegalStateException("Transaction " + transaction +
- " is not in a valid state to be invoking cache operations on.");
- TransactionXaAdapter current = tx2XaCacheMapping.get(transaction);
- if (current == null) {
- GlobalTransaction tx = rpcManager == null ? new GlobalTransaction(false) : new GlobalTransaction(rpcManager.getLocalAddress(), false);
- current = new TransactionXaAdapter(tx, icc, invoker, commandsFactory, configuration, this, transaction);
- tx2XaCacheMapping.put(transaction, current);
- transaction.enlistResource(current);
- notifier.notifyTransactionRegistered(tx, ctx);
- }
- return current;
- }
-
- public void delist(Transaction transaction) {
- if (transaction == null) throw new IllegalArgumentException("Null not allowed");
- TransactionXaAdapter xaAdapter = tx2XaCacheMapping.remove(transaction);
- if (xaAdapter == null) {
- throw new IllegalStateException("This method should only be called by a thread that has a tx association.");
- }
- }
-
- public boolean inTxScope() {
- try {
- if (tm == null) return false;
- Transaction transaction = tm.getTransaction();
- return transaction != null;
- } catch (SystemException e) {
- throw new CacheException(e);
- }
- }
-
- public Transaction getOngoingTx() {
- try {
- return tm.getTransaction();
- } catch (SystemException e) {
- throw new CacheException(e);
- }
- }
-
- public int getNumberOfInitiatedTx() {
- return tx2XaCacheMapping.size();
- }
-
- private boolean isValid(Transaction tx) {
- return isActive(tx) || isPreparing(tx);
- }
-
- /**
- * Returns true if transaction is PREPARING, false otherwise
- */
- private boolean isPreparing(Transaction tx) {
- if (tx == null) return false;
- int status;
- try {
- status = tx.getStatus();
- return status == Status.STATUS_PREPARING;
- }
- catch (SystemException e) {
- return false;
- }
- }
-
- /**
- * Returns true if transaction is ACTIVE, false otherwise
- */
- private boolean isActive(Transaction tx) {
- if (tx == null) return false;
- int status;
- try {
- status = tx.getStatus();
- return status == Status.STATUS_ACTIVE;
- }
- catch (SystemException e) {
- return false;
- }
- }
-
- public int getActiveLocallyInitiatedTxCount() {
- if (this.tx2XaCacheMapping == null) return 0;
- return tx2XaCacheMapping.size();
- }
-
- public int getActiveRemotelyInitiatedTxCount() {
- if (this.remoteTxMap == null) return 0;
- return this.remoteTxMap.size();
- }
-
- public TransactionXaAdapter getXaCache(Transaction tx) {
- return tx2XaCacheMapping.get(tx);
- }
-
- public Transaction getRunningTx() {
- try {
- return tm == null ? null : tm.getTransaction();
- } catch (SystemException e) {
- throw new CacheException(e);
- }
- }
-}
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-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/util/concurrent/locks/LockManagerImpl.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -25,7 +25,7 @@
import org.infinispan.container.entries.CacheEntry;
import org.infinispan.context.Flag;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.factories.annotations.Inject;
import org.infinispan.factories.annotations.Start;
import org.infinispan.jmx.annotations.MBean;
Modified: trunk/core/src/main/java/org/infinispan/util/concurrent/locks/OwnableReentrantLock.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/util/concurrent/locks/OwnableReentrantLock.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/util/concurrent/locks/OwnableReentrantLock.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -23,7 +23,7 @@
import net.jcip.annotations.ThreadSafe;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.util.logging.Log;
import org.infinispan.util.logging.LogFactory;
@@ -33,7 +33,7 @@
/**
* A lock that supports reentrancy based on owner (and not on current thread). For this to work, the lock needs to be
- * constructed with a reference to the {@link org.infinispan.context.container.InvocationContextContainer}, so it is able to determine whether the
+ * constructed with a reference to the {@link org.infinispan.context.InvocationContextContainer}, so it is able to determine whether the
* caller's "owner" reference is the current thread or a {@link org.infinispan.transaction.xa.GlobalTransaction} instance.
* <p/>
* This makes this lock implementation very closely tied to Infinispan internals, but it provides for a very clean,
Modified: trunk/core/src/main/java/org/infinispan/util/concurrent/locks/containers/OwnableReentrantPerEntryLockContainer.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/util/concurrent/locks/containers/OwnableReentrantPerEntryLockContainer.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/util/concurrent/locks/containers/OwnableReentrantPerEntryLockContainer.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -1,6 +1,6 @@
package org.infinispan.util.concurrent.locks.containers;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.util.concurrent.locks.OwnableReentrantLock;
import java.util.concurrent.locks.Lock;
Modified: trunk/core/src/main/java/org/infinispan/util/concurrent/locks/containers/OwnableReentrantStripedLockContainer.java
===================================================================
--- trunk/core/src/main/java/org/infinispan/util/concurrent/locks/containers/OwnableReentrantStripedLockContainer.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/main/java/org/infinispan/util/concurrent/locks/containers/OwnableReentrantStripedLockContainer.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -22,7 +22,7 @@
package org.infinispan.util.concurrent.locks.containers;
import net.jcip.annotations.ThreadSafe;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.util.concurrent.locks.OwnableReentrantLock;
import java.util.Arrays;
Modified: trunk/core/src/test/java/org/infinispan/api/mvcc/LockAssert.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/api/mvcc/LockAssert.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/test/java/org/infinispan/api/mvcc/LockAssert.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -1,7 +1,7 @@
package org.infinispan.api.mvcc;
import org.infinispan.Cache;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.context.InvocationContext;
import org.infinispan.context.impl.TxInvocationContext;
import org.infinispan.test.TestingUtil;
@@ -21,7 +21,7 @@
public static void assertNotLocked(Object key, InvocationContextContainer icc) {
// can't rely on the negative test since other entries may share the same lock with lock striping.
- assert !icc.getLocalInvocationContext(true).hasLockedKey(key) : key + " lock recorded!";
+ assert !icc.getLocalInvocationContext().hasLockedKey(key) : key + " lock recorded!";
}
public static void assertNoLocks(LockManager lockManager, InvocationContextContainer icc) {
Modified: trunk/core/src/test/java/org/infinispan/api/mvcc/LockPerEntryTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/api/mvcc/LockPerEntryTest.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/test/java/org/infinispan/api/mvcc/LockPerEntryTest.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -2,7 +2,7 @@
import org.infinispan.Cache;
import org.infinispan.config.Configuration;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.manager.CacheManager;
import org.infinispan.test.SingleCacheManagerTest;
import org.infinispan.test.TestingUtil;
Modified: trunk/core/src/test/java/org/infinispan/api/mvcc/LockTestBase.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/api/mvcc/LockTestBase.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/test/java/org/infinispan/api/mvcc/LockTestBase.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -2,7 +2,7 @@
import org.infinispan.Cache;
import org.infinispan.config.Configuration;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.manager.CacheManager;
import org.infinispan.test.TestingUtil;
import org.infinispan.test.fwk.TestCacheManagerFactory;
Modified: trunk/core/src/test/java/org/infinispan/api/mvcc/PutForExternalReadTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/api/mvcc/PutForExternalReadTest.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/test/java/org/infinispan/api/mvcc/PutForExternalReadTest.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -18,7 +18,7 @@
import org.infinispan.test.ReplListener;
import org.infinispan.test.TestingUtil;
import org.infinispan.transaction.lookup.DummyTransactionManagerLookup;
-import org.infinispan.transaction.xa.TxEnlistingManager;
+import org.infinispan.transaction.xa.TransactionTable;
import static org.testng.AssertJUnit.*;
import org.testng.annotations.Test;
@@ -226,10 +226,6 @@
cacheModeLocalTest(true);
}
- private TxEnlistingManager getTxEnlistingManager(Cache cache) {
- return TestingUtil.extractComponent(cache, TxEnlistingManager.class);
- }
-
/**
* Tests that suspended transactions do not leak. See JBCACHE-1246.
*/
@@ -240,13 +236,13 @@
tm1.commit();
replListener2.waitForRpc();
- TxEnlistingManager tt1 = getTxEnlistingManager(cache1);
- TxEnlistingManager tt2 = getTxEnlistingManager(cache2);
+ TransactionTable tt1 = TestingUtil.extractComponent(cache1, TransactionTable.class);
+ TransactionTable tt2 = TestingUtil.extractComponent(cache2, TransactionTable.class);
- assert tt1.getActiveRemotelyInitiatedTxCount() == 0 : "Cache 1 should have no stale global TXs";
- assert tt1.getActiveLocallyInitiatedTxCount() == 0 : "Cache 1 should have no stale local TXs";
- assert tt2.getActiveRemotelyInitiatedTxCount() == 0 : "Cache 2 should have no stale global TXs";
- assert tt2.getActiveLocallyInitiatedTxCount() == 0 : "Cache 2 should have no stale local TXs";
+ assert tt1.getRemoteTxCount() == 0 : "Cache 1 should have no stale global TXs";
+ assert tt1.getLocalTxCount() == 0 : "Cache 1 should have no stale local TXs";
+ assert tt2.getRemoteTxCount() == 0 : "Cache 2 should have no stale global TXs";
+ assert tt2.getLocalTxCount() == 0 : "Cache 2 should have no stale local TXs";
System.out.println("PutForExternalReadTest.testMemLeakOnSuspendedTransactions");
replListener2.expectWithTx(PutKeyValueCommand.class);
@@ -256,10 +252,10 @@
tm1.commit();
replListener2.waitForRpc();
- assert tt1.getActiveRemotelyInitiatedTxCount() == 0 : "Cache 1 should have no stale global TXs";
- assert tt1.getActiveLocallyInitiatedTxCount() == 0 : "Cache 1 should have no stale local TXs";
- assert tt2.getActiveRemotelyInitiatedTxCount() == 0 : "Cache 2 should have no stale global TXs";
- assert tt2.getActiveLocallyInitiatedTxCount() == 0 : "Cache 2 should have no stale local TXs";
+ assert tt1.getRemoteTxCount() == 0 : "Cache 1 should have no stale global TXs";
+ assert tt1.getLocalTxCount() == 0 : "Cache 1 should have no stale local TXs";
+ assert tt2.getRemoteTxCount() == 0 : "Cache 2 should have no stale global TXs";
+ assert tt2.getLocalTxCount() == 0 : "Cache 2 should have no stale local TXs";
replListener2.expectWithTx(PutKeyValueCommand.class);
tm1.begin();
@@ -268,10 +264,10 @@
tm1.commit();
replListener2.waitForRpc();
- assert tt1.getActiveRemotelyInitiatedTxCount() == 0 : "Cache 1 should have no stale global TXs";
- assert tt1.getActiveLocallyInitiatedTxCount() == 0 : "Cache 1 should have no stale local TXs";
- assert tt2.getActiveRemotelyInitiatedTxCount() == 0 : "Cache 2 should have no stale global TXs";
- assert tt2.getActiveLocallyInitiatedTxCount() == 0 : "Cache 2 should have no stale local TXs";
+ assert tt1.getRemoteTxCount() == 0 : "Cache 1 should have no stale global TXs";
+ assert tt1.getLocalTxCount() == 0 : "Cache 1 should have no stale local TXs";
+ assert tt2.getRemoteTxCount() == 0 : "Cache 2 should have no stale global TXs";
+ assert tt2.getLocalTxCount() == 0 : "Cache 2 should have no stale local TXs";
replListener2.expectWithTx(PutKeyValueCommand.class, PutKeyValueCommand.class);
tm1.begin();
@@ -281,10 +277,10 @@
tm1.commit();
replListener2.waitForRpc();
- assert tt1.getActiveRemotelyInitiatedTxCount() == 0 : "Cache 1 should have no stale global TXs";
- assert tt1.getActiveLocallyInitiatedTxCount() == 0 : "Cache 1 should have no stale local TXs";
- assert tt2.getActiveRemotelyInitiatedTxCount() == 0 : "Cache 2 should have no stale global TXs";
- assert tt2.getActiveLocallyInitiatedTxCount() == 0 : "Cache 2 should have no stale local TXs";
+ assert tt1.getRemoteTxCount() == 0 : "Cache 1 should have no stale global TXs";
+ assert tt1.getLocalTxCount() == 0 : "Cache 1 should have no stale local TXs";
+ assert tt2.getRemoteTxCount() == 0 : "Cache 2 should have no stale global TXs";
+ assert tt2.getLocalTxCount() == 0 : "Cache 2 should have no stale local TXs";
}
/**
Modified: trunk/core/src/test/java/org/infinispan/api/mvcc/repeatable_read/WriteSkewTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/api/mvcc/repeatable_read/WriteSkewTest.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/test/java/org/infinispan/api/mvcc/repeatable_read/WriteSkewTest.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -3,7 +3,7 @@
import org.infinispan.Cache;
import org.infinispan.api.mvcc.LockAssert;
import org.infinispan.config.Configuration;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.manager.CacheManager;
import org.infinispan.test.TestingUtil;
import org.infinispan.test.fwk.TestCacheManagerFactory;
Modified: trunk/core/src/test/java/org/infinispan/atomic/AtomicMapFunctionalTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/atomic/AtomicMapFunctionalTest.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/test/java/org/infinispan/atomic/AtomicMapFunctionalTest.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -4,7 +4,7 @@
import org.infinispan.config.Configuration;
import static org.infinispan.context.Flag.SKIP_LOCKING;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.manager.CacheManager;
import org.infinispan.test.TestingUtil;
import org.infinispan.test.fwk.TestCacheManagerFactory;
@@ -73,10 +73,10 @@
AtomicMap<String, String> map = cache.getAtomicMap("key");
assert map.isEmpty();
InvocationContextContainer icc = TestingUtil.extractComponent(cache, InvocationContextContainer.class);
- InvocationContext ic = icc.getLocalInvocationContext(true);
+ InvocationContext ic = icc.getLocalInvocationContext();
ic.setFlags(SKIP_LOCKING);
log.debug("Doing a put");
- assert icc.getLocalInvocationContext(false).hasFlag(SKIP_LOCKING);
+ assert icc.getThreadContext().hasFlag(SKIP_LOCKING);
map.put("a", "b");
log.debug("Put complete");
assert map.get("a").equals("b");
@@ -89,7 +89,7 @@
AtomicMap<String, String> map = cache.getAtomicMap("key");
tm.begin();
assert map.isEmpty();
- TestingUtil.extractComponent(cache, InvocationContextContainer.class).getLocalInvocationContext(true).setFlags(SKIP_LOCKING);
+ TestingUtil.extractComponent(cache, InvocationContextContainer.class).getLocalInvocationContext().setFlags(SKIP_LOCKING);
map.put("a", "b");
assert map.get("a").equals("b");
Transaction t = tm.suspend();
@@ -108,7 +108,7 @@
assert map.isEmpty();
map.put("x", "y");
assert map.get("x").equals("y");
- TestingUtil.extractComponent(cache, InvocationContextContainer.class).getLocalInvocationContext(true).setFlags(SKIP_LOCKING);
+ TestingUtil.extractComponent(cache, InvocationContextContainer.class).getLocalInvocationContext().setFlags(SKIP_LOCKING);
log.debug("Doing a put");
map.put("a", "b");
log.debug("Put complete");
Modified: trunk/core/src/test/java/org/infinispan/notifications/cachelistener/CacheNotifierImplTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/notifications/cachelistener/CacheNotifierImplTest.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/test/java/org/infinispan/notifications/cachelistener/CacheNotifierImplTest.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -4,8 +4,8 @@
import static org.easymock.classextension.EasyMock.createNiceMock;
import org.infinispan.Cache;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
-import org.infinispan.context.container.ReplicationInvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainerImpl;
import org.infinispan.context.impl.NonTxInvocationContext;
import org.infinispan.notifications.cachelistener.event.CacheEntryEvent;
import org.infinispan.notifications.cachelistener.event.CacheEntryModifiedEvent;
@@ -29,7 +29,7 @@
n = new CacheNotifierImpl();
mockCache = createNiceMock(Cache.class);
EasyMock.replay(mockCache);
- InvocationContextContainer icc = new ReplicationInvocationContextContainer();
+ InvocationContextContainer icc = new InvocationContextContainerImpl();
n.injectDependencies(icc, mockCache);
cl = new CacheListener();
n.start();
Modified: trunk/core/src/test/java/org/infinispan/test/AbstractCacheTest.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/test/AbstractCacheTest.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/test/java/org/infinispan/test/AbstractCacheTest.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -46,7 +46,7 @@
removeInMemoryData(cache);
clearCacheLoader(cache);
clearReplicationQueues(cache);
- InvocationContext invocationContext = ((AdvancedCache) cache).getInvocationContextContainer().getLocalInvocationContext(true);
+ InvocationContext invocationContext = ((AdvancedCache) cache).getInvocationContextContainer().getLocalInvocationContext();
}
}
}
Modified: trunk/core/src/test/java/org/infinispan/test/TestingUtil.java
===================================================================
--- trunk/core/src/test/java/org/infinispan/test/TestingUtil.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/core/src/test/java/org/infinispan/test/TestingUtil.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -13,7 +13,7 @@
import org.infinispan.commands.CommandsFactory;
import org.infinispan.commands.VisitableCommand;
import org.infinispan.context.InvocationContext;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.factories.ComponentRegistry;
import org.infinispan.factories.GlobalComponentRegistry;
import org.infinispan.interceptors.InterceptorChain;
@@ -523,7 +523,7 @@
ComponentRegistry cr = extractComponentRegistry(cache);
InterceptorChain ic = cr.getComponent(InterceptorChain.class);
InvocationContextContainer icc = cr.getComponent(InvocationContextContainer.class);
- InvocationContext ctxt = icc.getLocalInvocationContext(true);
+ InvocationContext ctxt = icc.getLocalInvocationContext();
ic.invoke(ctxt, command);
}
Modified: trunk/tree/src/main/java/org/infinispan/tree/NodeImpl.java
===================================================================
--- trunk/tree/src/main/java/org/infinispan/tree/NodeImpl.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/tree/src/main/java/org/infinispan/tree/NodeImpl.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -26,7 +26,7 @@
import org.infinispan.atomic.atomichashmap.AtomicHashMapProxy;
import org.infinispan.batch.BatchContainer;
import org.infinispan.context.Flag;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.util.Immutables;
import org.infinispan.util.Util;
@@ -59,7 +59,7 @@
}
public Node<K, V> getParent(Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return getParent();
}
@@ -79,7 +79,7 @@
}
public Set<Node<K, V>> getChildren(Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return getChildren();
}
@@ -88,7 +88,7 @@
}
public Set<Object> getChildrenNames(Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return getChildrenNames();
}
@@ -98,7 +98,7 @@
}
public Map<K, V> getData(Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return getData();
}
@@ -113,7 +113,7 @@
}
public Set<K> getKeys(Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return getKeys();
}
@@ -141,7 +141,7 @@
}
public Node<K, V> addChild(Fqn f, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return addChild(f);
}
@@ -150,7 +150,7 @@
}
public boolean removeChild(Fqn f, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return removeChild(f);
}
@@ -176,7 +176,7 @@
}
public boolean removeChild(Object childName, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return removeChild(childName);
}
@@ -194,7 +194,7 @@
}
public Node<K, V> getChild(Fqn f, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return getChild(f);
}
@@ -212,7 +212,7 @@
}
public Node<K, V> getChild(Object name, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return getChild(name);
}
@@ -228,7 +228,7 @@
}
public V put(K key, V value, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return put(key, value);
}
@@ -246,7 +246,7 @@
}
public V putIfAbsent(K key, V value, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return putIfAbsent(key, value);
}
@@ -265,7 +265,7 @@
}
public V replace(K key, V value, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return replace(key, value);
}
@@ -286,7 +286,7 @@
}
public boolean replace(K key, V oldValue, V value, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return replace(key, oldValue, value);
}
@@ -301,7 +301,7 @@
}
public void putAll(Map<? extends K, ? extends V> map, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
putAll(map);
}
@@ -318,7 +318,7 @@
}
public void replaceAll(Map<? extends K, ? extends V> map, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
replaceAll(map);
}
@@ -327,7 +327,7 @@
}
public V get(K key, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return get(key);
}
@@ -342,7 +342,7 @@
}
public V remove(K key, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return remove(key);
}
@@ -351,7 +351,7 @@
}
public void clearData(Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
clearData();
}
@@ -360,7 +360,7 @@
}
public int dataSize(Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return dataSize();
}
@@ -375,7 +375,7 @@
}
public boolean hasChild(Fqn f, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return hasChild(f);
}
@@ -384,7 +384,7 @@
}
public boolean hasChild(Object o, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return hasChild(o);
}
@@ -404,7 +404,7 @@
}
public void removeChildren(Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
removeChildren();
}
Modified: trunk/tree/src/main/java/org/infinispan/tree/TreeCacheImpl.java
===================================================================
--- trunk/tree/src/main/java/org/infinispan/tree/TreeCacheImpl.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/tree/src/main/java/org/infinispan/tree/TreeCacheImpl.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -52,7 +52,7 @@
}
public Node<K, V> getRoot(Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return getRoot();
}
@@ -61,7 +61,7 @@
}
public V put(String fqn, K key, V value, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return put(fqn, key, value);
}
@@ -76,7 +76,7 @@
}
public void put(Fqn fqn, Map<? extends K, ? extends V> data, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
put(fqn, data);
}
@@ -85,7 +85,7 @@
}
public void put(String fqn, Map<? extends K, ? extends V> data, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
put(fqn, data);
}
@@ -102,7 +102,7 @@
}
public V remove(Fqn fqn, K key, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return remove(fqn, key);
}
@@ -111,7 +111,7 @@
}
public V remove(String fqn, K key, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return remove(fqn, key);
}
@@ -132,7 +132,7 @@
}
public boolean removeNode(Fqn fqn, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return removeNode(fqn);
}
@@ -141,7 +141,7 @@
}
public boolean removeNode(String fqn, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return removeNode(fqn);
}
@@ -167,7 +167,7 @@
}
public Node<K, V> getNode(String fqn, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return getNode(fqn);
}
@@ -179,7 +179,7 @@
}
public V get(Fqn fqn, K key, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return get(fqn, key);
}
@@ -188,12 +188,12 @@
}
public boolean exists(String fqn, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return exists(fqn);
}
public boolean exists(Fqn fqn, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return exists(fqn);
}
@@ -202,7 +202,7 @@
}
public V get(String fqn, K key, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return get(fqn, key);
}
@@ -251,7 +251,7 @@
}
public void move(Fqn nodeToMove, Fqn newParent, Flag... flags) throws NodeNotExistsException {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
move(nodeToMove, newParent);
}
@@ -260,7 +260,7 @@
}
public void move(String nodeToMove, String newParent, Flag... flags) throws NodeNotExistsException {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
move(nodeToMove, newParent);
}
@@ -275,7 +275,7 @@
}
public Map<K, V> getData(Fqn fqn, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return getData(fqn);
}
@@ -284,7 +284,7 @@
}
public Set<K> getKeys(String fqn, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return getKeys(fqn);
}
@@ -299,7 +299,7 @@
}
public Set<K> getKeys(Fqn fqn, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return getKeys(fqn);
}
@@ -308,7 +308,7 @@
}
public void clearData(String fqn, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
}
public void clearData(Fqn fqn) {
@@ -322,7 +322,7 @@
}
public void clearData(Fqn fqn, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
}
@SuppressWarnings("unchecked")
@@ -340,7 +340,7 @@
}
public V put(Fqn fqn, K key, V value, Flag... flags) {
- icc.getLocalInvocationContext(true).setFlags(flags);
+ icc.getLocalInvocationContext().setFlags(flags);
return put(fqn, key, value);
}
Modified: trunk/tree/src/main/java/org/infinispan/tree/TreeStructureSupport.java
===================================================================
--- trunk/tree/src/main/java/org/infinispan/tree/TreeStructureSupport.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/tree/src/main/java/org/infinispan/tree/TreeStructureSupport.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -27,7 +27,7 @@
import org.infinispan.batch.AutoBatchSupport;
import org.infinispan.batch.BatchContainer;
import org.infinispan.context.Flag;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.util.concurrent.locks.LockManager;
import org.infinispan.util.logging.Log;
import org.infinispan.util.logging.LogFactory;
Modified: trunk/tree/src/test/java/org/infinispan/api/tree/NodeMoveAPITest.java
===================================================================
--- trunk/tree/src/test/java/org/infinispan/api/tree/NodeMoveAPITest.java 2009-05-12 20:30:27 UTC (rev 266)
+++ trunk/tree/src/test/java/org/infinispan/api/tree/NodeMoveAPITest.java 2009-05-13 03:15:11 UTC (rev 267)
@@ -3,7 +3,7 @@
import org.infinispan.api.mvcc.LockAssert;
import org.infinispan.config.Configuration;
import org.infinispan.container.DataContainer;
-import org.infinispan.context.container.InvocationContextContainer;
+import org.infinispan.context.InvocationContextContainer;
import org.infinispan.factories.ComponentRegistry;
import org.infinispan.manager.CacheManager;
import org.infinispan.test.SingleCacheManagerTest;
More information about the infinispan-commits
mailing list