X-Git-Url: https://git.opendaylight.org/gerrit/gitweb?a=blobdiff_plain;ds=sidebyside;f=opendaylight%2Fmd-sal%2Fsal-distributed-datastore%2Fsrc%2Fmain%2Fjava%2Forg%2Fopendaylight%2Fcontroller%2Fcluster%2Fdatastore%2FTransactionProxy.java;h=667b64ce18839386f0c9ce44375bd7f88649588a;hb=f97618f25dfc073d1de5d883f1794eefdb3e5c16;hp=d79cd6f69f4b0e3e4f1171a332b45c36adbbd515;hpb=4a2e88ee3ff4d07a980083586328d3a52639e204;p=controller.git
diff --git a/opendaylight/md-sal/sal-distributed-datastore/src/main/java/org/opendaylight/controller/cluster/datastore/TransactionProxy.java b/opendaylight/md-sal/sal-distributed-datastore/src/main/java/org/opendaylight/controller/cluster/datastore/TransactionProxy.java
index d79cd6f69f..667b64ce18 100644
--- a/opendaylight/md-sal/sal-distributed-datastore/src/main/java/org/opendaylight/controller/cluster/datastore/TransactionProxy.java
+++ b/opendaylight/md-sal/sal-distributed-datastore/src/main/java/org/opendaylight/controller/cluster/datastore/TransactionProxy.java
@@ -18,11 +18,11 @@ import com.google.common.base.Optional;
import com.google.common.base.Preconditions;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.CheckedFuture;
-import com.google.common.util.concurrent.FutureCallback;
-import com.google.common.util.concurrent.Futures;
+import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.SettableFuture;
import java.util.ArrayList;
import java.util.Collection;
+import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -32,6 +32,7 @@ import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicLong;
import javax.annotation.concurrent.GuardedBy;
+import org.opendaylight.controller.cluster.datastore.compat.PreLithiumTransactionContextImpl;
import org.opendaylight.controller.cluster.datastore.exceptions.NoShardLeaderException;
import org.opendaylight.controller.cluster.datastore.identifiers.TransactionIdentifier;
import org.opendaylight.controller.cluster.datastore.messages.CloseTransaction;
@@ -40,6 +41,7 @@ import org.opendaylight.controller.cluster.datastore.messages.CreateTransactionR
import org.opendaylight.controller.cluster.datastore.shardstrategy.ShardStrategyFactory;
import org.opendaylight.controller.cluster.datastore.utils.ActorContext;
import org.opendaylight.controller.md.sal.common.api.data.ReadFailedException;
+import org.opendaylight.controller.sal.core.spi.data.AbstractDOMStoreTransaction;
import org.opendaylight.controller.sal.core.spi.data.DOMStoreReadWriteTransaction;
import org.opendaylight.controller.sal.core.spi.data.DOMStoreThreePhaseCommitCohort;
import org.opendaylight.yangtools.util.concurrent.MappingCheckedFuture;
@@ -64,12 +66,23 @@ import scala.concurrent.duration.FiniteDuration;
* shards will be executed.
*
*/
-public class TransactionProxy implements DOMStoreReadWriteTransaction {
+public class TransactionProxy extends AbstractDOMStoreTransaction implements DOMStoreReadWriteTransaction {
public static enum TransactionType {
READ_ONLY,
WRITE_ONLY,
- READ_WRITE
+ READ_WRITE;
+
+ // Cache all values
+ private static final TransactionType[] VALUES = values();
+
+ public static TransactionType fromInt(final int type) {
+ try {
+ return VALUES[type];
+ } catch (IndexOutOfBoundsException e) {
+ throw new IllegalArgumentException("In TransactionType enum value " + type, e);
+ }
+ }
}
static final Mapper SAME_FAILURE_TRANSFORMER =
@@ -134,7 +147,7 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
remoteTransactionActors = referent.remoteTransactionActors;
remoteTransactionActorsMB = referent.remoteTransactionActorsMB;
actorContext = referent.actorContext;
- identifier = referent.identifier;
+ identifier = referent.getIdentifier();
}
@Override
@@ -164,7 +177,7 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
* PhantomReference.
*/
private List remoteTransactionActors;
- private AtomicBoolean remoteTransactionActorsMB;
+ private volatile AtomicBoolean remoteTransactionActorsMB;
/**
* Stores the create transaction results per shard.
@@ -173,19 +186,20 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
private final TransactionType transactionType;
private final ActorContext actorContext;
- private final TransactionIdentifier identifier;
private final String transactionChainId;
private final SchemaContext schemaContext;
private boolean inReadyState;
- private final Semaphore operationLimiter;
- private final OperationCompleter operationCompleter;
+
+ private volatile boolean initialized;
+ private Semaphore operationLimiter;
+ private OperationCompleter operationCompleter;
public TransactionProxy(ActorContext actorContext, TransactionType transactionType) {
this(actorContext, transactionType, "");
}
- public TransactionProxy(ActorContext actorContext, TransactionType transactionType,
- String transactionChainId) {
+ public TransactionProxy(ActorContext actorContext, TransactionType transactionType, String transactionChainId) {
+ super(createIdentifier(actorContext));
this.actorContext = Preconditions.checkNotNull(actorContext,
"actorContext should not be null");
this.transactionType = Preconditions.checkNotNull(transactionType,
@@ -194,32 +208,16 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
"schemaContext should not be null");
this.transactionChainId = transactionChainId;
+ LOG.debug("Created txn {} of type {} on chain {}", getIdentifier(), transactionType, transactionChainId);
+ }
+
+ private static TransactionIdentifier createIdentifier(ActorContext actorContext) {
String memberName = actorContext.getCurrentMemberName();
- if(memberName == null){
+ if (memberName == null) {
memberName = "UNKNOWN-MEMBER";
}
- this.identifier = TransactionIdentifier.builder().memberName(memberName).counter(
- counter.getAndIncrement()).build();
-
- if(transactionType == TransactionType.READ_ONLY) {
- // Read-only Tx's aren't explicitly closed by the client so we create a PhantomReference
- // to close the remote Tx's when this instance is no longer in use and is garbage
- // collected.
-
- remoteTransactionActors = Lists.newArrayList();
- remoteTransactionActorsMB = new AtomicBoolean();
-
- TransactionProxyCleanupPhantomReference cleanup =
- new TransactionProxyCleanupPhantomReference(this);
- phantomReferenceCache.put(cleanup, cleanup);
- }
-
- // Note : Currently mailbox-capacity comes from akka.conf and not from the config-subsystem
- this.operationLimiter = new Semaphore(actorContext.getTransactionOutstandingOperationLimit());
- this.operationCompleter = new OperationCompleter(operationLimiter);
-
- LOG.debug("Created txn {} of type {} on chain {}", identifier, transactionType, transactionChainId);
+ return new TransactionIdentifier(memberName, counter.getAndIncrement());
}
@VisibleForTesting
@@ -253,18 +251,21 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
Preconditions.checkState(transactionType != TransactionType.WRITE_ONLY,
"Read operation on write-only transaction is not allowed");
- LOG.debug("Tx {} read {}", identifier, path);
+ LOG.debug("Tx {} read {}", getIdentifier(), path);
throttleOperation();
+ final SettableFuture>> proxyFuture = SettableFuture.create();
+
TransactionFutureCallback txFutureCallback = getOrCreateTxFutureCallback(path);
- return txFutureCallback.enqueueReadOperation(new ReadOperation>>() {
+ txFutureCallback.enqueueTransactionOperation(new TransactionOperation() {
@Override
- public CheckedFuture>, ReadFailedException> invoke(
- TransactionContext transactionContext) {
- return transactionContext.readData(path);
+ public void invoke(TransactionContext transactionContext) {
+ transactionContext.readData(path, proxyFuture);
}
});
+
+ return MappingCheckedFuture.create(proxyFuture, ReadFailedException.MAPPER);
}
@Override
@@ -273,19 +274,22 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
Preconditions.checkState(transactionType != TransactionType.WRITE_ONLY,
"Exists operation on write-only transaction is not allowed");
- LOG.debug("Tx {} exists {}", identifier, path);
+ LOG.debug("Tx {} exists {}", getIdentifier(), path);
throttleOperation();
+ final SettableFuture proxyFuture = SettableFuture.create();
+
TransactionFutureCallback txFutureCallback = getOrCreateTxFutureCallback(path);
- return txFutureCallback.enqueueReadOperation(new ReadOperation() {
+ txFutureCallback.enqueueTransactionOperation(new TransactionOperation() {
@Override
- public CheckedFuture invoke(TransactionContext transactionContext) {
- return transactionContext.dataExists(path);
+ public void invoke(TransactionContext transactionContext) {
+ transactionContext.dataExists(path, proxyFuture);
}
});
- }
+ return MappingCheckedFuture.create(proxyFuture, ReadFailedException.MAPPER);
+ }
private void checkModificationState() {
Preconditions.checkState(transactionType != TransactionType.READ_ONLY,
@@ -299,8 +303,19 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
}
private void throttleOperation(int acquirePermits) {
+ if(!initialized) {
+ // Note : Currently mailbox-capacity comes from akka.conf and not from the config-subsystem
+ operationLimiter = new Semaphore(actorContext.getTransactionOutstandingOperationLimit());
+ operationCompleter = new OperationCompleter(operationLimiter);
+
+ // Make sure we write this last because it's volatile and will also publish the non-volatile writes
+ // above as well so they'll be visible to other threads.
+ initialized = true;
+ }
+
try {
- if(!operationLimiter.tryAcquire(acquirePermits, actorContext.getDatastoreContext().getOperationTimeoutInSeconds(), TimeUnit.SECONDS)){
+ if(!operationLimiter.tryAcquire(acquirePermits,
+ actorContext.getDatastoreContext().getOperationTimeoutInSeconds(), TimeUnit.SECONDS)){
LOG.warn("Failed to acquire operation permit for transaction {}", getIdentifier());
}
} catch (InterruptedException e) {
@@ -318,12 +333,12 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
checkModificationState();
- LOG.debug("Tx {} write {}", identifier, path);
+ LOG.debug("Tx {} write {}", getIdentifier(), path);
throttleOperation();
TransactionFutureCallback txFutureCallback = getOrCreateTxFutureCallback(path);
- txFutureCallback.enqueueModifyOperation(new TransactionOperation() {
+ txFutureCallback.enqueueTransactionOperation(new TransactionOperation() {
@Override
public void invoke(TransactionContext transactionContext) {
transactionContext.writeData(path, data);
@@ -336,12 +351,12 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
checkModificationState();
- LOG.debug("Tx {} merge {}", identifier, path);
+ LOG.debug("Tx {} merge {}", getIdentifier(), path);
throttleOperation();
TransactionFutureCallback txFutureCallback = getOrCreateTxFutureCallback(path);
- txFutureCallback.enqueueModifyOperation(new TransactionOperation() {
+ txFutureCallback.enqueueTransactionOperation(new TransactionOperation() {
@Override
public void invoke(TransactionContext transactionContext) {
transactionContext.mergeData(path, data);
@@ -354,12 +369,12 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
checkModificationState();
- LOG.debug("Tx {} delete {}", identifier, path);
+ LOG.debug("Tx {} delete {}", getIdentifier(), path);
throttleOperation();
TransactionFutureCallback txFutureCallback = getOrCreateTxFutureCallback(path);
- txFutureCallback.enqueueModifyOperation(new TransactionOperation() {
+ txFutureCallback.enqueueTransactionOperation(new TransactionOperation() {
@Override
public void invoke(TransactionContext transactionContext) {
transactionContext.deleteData(path);
@@ -372,26 +387,41 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
checkModificationState();
- throttleOperation(txFutureCallbackMap.size());
-
inReadyState = true;
- LOG.debug("Tx {} Readying {} transactions for commit", identifier,
+ LOG.debug("Tx {} Readying {} transactions for commit", getIdentifier(),
txFutureCallbackMap.size());
+ if(txFutureCallbackMap.size() == 0) {
+ onTransactionReady(Collections.>emptyList());
+ TransactionRateLimitingCallback.adjustRateLimitForUnusedTransaction(actorContext);
+ return NoOpDOMStoreThreePhaseCommitCohort.INSTANCE;
+ }
+
+ throttleOperation(txFutureCallbackMap.size());
+
List> cohortFutures = Lists.newArrayList();
for(TransactionFutureCallback txFutureCallback : txFutureCallbackMap.values()) {
- LOG.debug("Tx {} Readying transaction for shard {} chain {}", identifier,
+ LOG.debug("Tx {} Readying transaction for shard {} chain {}", getIdentifier(),
txFutureCallback.getShardName(), transactionChainId);
- Future future = txFutureCallback.enqueueFutureOperation(new FutureOperation() {
- @Override
- public Future invoke(TransactionContext transactionContext) {
- return transactionContext.readyTransaction();
- }
- });
+ final TransactionContext transactionContext = txFutureCallback.getTransactionContext();
+ final Future future;
+ if (transactionContext != null) {
+ // avoid the creation of a promise and a TransactionOperation
+ future = transactionContext.readyTransaction();
+ } else {
+ final Promise promise = akka.dispatch.Futures.promise();
+ txFutureCallback.enqueueTransactionOperation(new TransactionOperation() {
+ @Override
+ public void invoke(TransactionContext transactionContext) {
+ promise.completeWith(transactionContext.readyTransaction());
+ }
+ });
+ future = promise.future();
+ }
cohortFutures.add(future);
}
@@ -399,7 +429,7 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
onTransactionReady(cohortFutures);
return new ThreePhaseCommitCohortProxy(actorContext, cohortFutures,
- identifier.toString());
+ getIdentifier().toString());
}
/**
@@ -410,27 +440,10 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
protected void onTransactionReady(List> cohortFutures) {
}
- /**
- * Method called to send a CreateTransaction message to a shard.
- *
- * @param shard the shard actor to send to
- * @param serializedCreateMessage the serialized message to send
- * @return the response Future
- */
- protected Future