X-Git-Url: https://git.opendaylight.org/gerrit/gitweb?p=controller.git;a=blobdiff_plain;f=opendaylight%2Fmd-sal%2Fsal-distributed-datastore%2Fsrc%2Fmain%2Fjava%2Forg%2Fopendaylight%2Fcontroller%2Fcluster%2Fdatastore%2FTransactionProxy.java;h=7703f484c73e687391ff5dd9ffe1739afd201c0f;hp=ffb1ab7c55064ecf7683e1857800b256c9c82a1d;hb=112e6b1bfeed9c7125a073d1015d05e31f006bbf;hpb=f5a373c5378af41f62a2c36ced4046fbdb77e00b
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 ffb1ab7c55..7703f484c7 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
@@ -21,6 +21,14 @@ 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.SettableFuture;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+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.exceptions.NoShardLeaderException;
import org.opendaylight.controller.cluster.datastore.identifiers.TransactionIdentifier;
import org.opendaylight.controller.cluster.datastore.messages.CloseTransaction;
@@ -34,6 +42,7 @@ import org.opendaylight.controller.cluster.datastore.messages.ReadData;
import org.opendaylight.controller.cluster.datastore.messages.ReadDataReply;
import org.opendaylight.controller.cluster.datastore.messages.ReadyTransaction;
import org.opendaylight.controller.cluster.datastore.messages.ReadyTransactionReply;
+import org.opendaylight.controller.cluster.datastore.messages.SerializableMessage;
import org.opendaylight.controller.cluster.datastore.messages.WriteData;
import org.opendaylight.controller.cluster.datastore.shardstrategy.ShardStrategyFactory;
import org.opendaylight.controller.cluster.datastore.utils.ActorContext;
@@ -50,15 +59,6 @@ import scala.concurrent.Future;
import scala.concurrent.Promise;
import scala.concurrent.duration.FiniteDuration;
-import javax.annotation.concurrent.GuardedBy;
-import java.util.HashMap;
-import java.util.List;
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.TimeUnit;
-import java.util.concurrent.atomic.AtomicBoolean;
-import java.util.concurrent.atomic.AtomicLong;
-
/**
* TransactionProxy acts as a proxy for one or more transactions that were created on a remote shard
*
@@ -182,23 +182,23 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
private final TransactionType transactionType;
private final ActorContext actorContext;
private final TransactionIdentifier identifier;
- private final TransactionChainProxy transactionChainProxy;
+ private final String transactionChainId;
private final SchemaContext schemaContext;
private boolean inReadyState;
public TransactionProxy(ActorContext actorContext, TransactionType transactionType) {
- this(actorContext, transactionType, null);
+ this(actorContext, transactionType, "");
}
public TransactionProxy(ActorContext actorContext, TransactionType transactionType,
- TransactionChainProxy transactionChainProxy) {
+ String transactionChainId) {
this.actorContext = Preconditions.checkNotNull(actorContext,
"actorContext should not be null");
this.transactionType = Preconditions.checkNotNull(transactionType,
"transactionType should not be null");
this.schemaContext = Preconditions.checkNotNull(actorContext.getSchemaContext(),
"schemaContext should not be null");
- this.transactionChainProxy = transactionChainProxy;
+ this.transactionChainId = transactionChainId;
String memberName = actorContext.getCurrentMemberName();
if(memberName == null){
@@ -221,7 +221,7 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
phantomReferenceCache.put(cleanup, cleanup);
}
- LOG.debug("Created txn {} of type {}", identifier, transactionType);
+ LOG.debug("Created txn {} of type {} on chain {}", identifier, transactionType, transactionChainId);
}
@VisibleForTesting
@@ -237,9 +237,20 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
return recordedOperationFutures;
}
+ @VisibleForTesting
+ boolean hasTransactionContext() {
+ for(TransactionFutureCallback txFutureCallback : txFutureCallbackMap.values()) {
+ TransactionContext transactionContext = txFutureCallback.getTransactionContext();
+ if(transactionContext != null) {
+ return true;
+ }
+ }
+
+ return false;
+ }
+
@Override
- public CheckedFuture>, ReadFailedException> read(
- final YangInstanceIdentifier path) {
+ public CheckedFuture>, ReadFailedException> read(final YangInstanceIdentifier path) {
Preconditions.checkState(transactionType != TransactionType.WRITE_ONLY,
"Read operation on write-only transaction is not allowed");
@@ -247,37 +258,13 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
LOG.debug("Tx {} read {}", identifier, path);
TransactionFutureCallback txFutureCallback = getOrCreateTxFutureCallback(path);
- TransactionContext transactionContext = txFutureCallback.getTransactionContext();
-
- CheckedFuture>, ReadFailedException> future;
- if(transactionContext != null) {
- future = transactionContext.readData(path);
- } else {
- // The shard Tx hasn't been created yet so add the Tx operation to the Tx Future
- // callback to be executed after the Tx is created.
- final SettableFuture>> proxyFuture = SettableFuture.create();
- txFutureCallback.addTxOperationOnComplete(new TransactionOperation() {
- @Override
- public void invoke(TransactionContext transactionContext) {
- Futures.addCallback(transactionContext.readData(path),
- new FutureCallback>>() {
- @Override
- public void onSuccess(Optional> data) {
- proxyFuture.set(data);
- }
-
- @Override
- public void onFailure(Throwable t) {
- proxyFuture.setException(t);
- }
- });
- }
- });
-
- future = MappingCheckedFuture.create(proxyFuture, ReadFailedException.MAPPER);
- }
-
- return future;
+ return txFutureCallback.enqueueReadOperation(new ReadOperation>>() {
+ @Override
+ public CheckedFuture>, ReadFailedException> invoke(
+ TransactionContext transactionContext) {
+ return transactionContext.readData(path);
+ }
+ });
}
@Override
@@ -289,39 +276,15 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
LOG.debug("Tx {} exists {}", identifier, path);
TransactionFutureCallback txFutureCallback = getOrCreateTxFutureCallback(path);
- TransactionContext transactionContext = txFutureCallback.getTransactionContext();
-
- CheckedFuture future;
- if(transactionContext != null) {
- future = transactionContext.dataExists(path);
- } else {
- // The shard Tx hasn't been created yet so add the Tx operation to the Tx Future
- // callback to be executed after the Tx is created.
- final SettableFuture proxyFuture = SettableFuture.create();
- txFutureCallback.addTxOperationOnComplete(new TransactionOperation() {
- @Override
- public void invoke(TransactionContext transactionContext) {
- Futures.addCallback(transactionContext.dataExists(path),
- new FutureCallback() {
- @Override
- public void onSuccess(Boolean exists) {
- proxyFuture.set(exists);
- }
-
- @Override
- public void onFailure(Throwable t) {
- proxyFuture.setException(t);
- }
- });
- }
- });
-
- future = MappingCheckedFuture.create(proxyFuture, ReadFailedException.MAPPER);
- }
-
- return future;
+ return txFutureCallback.enqueueReadOperation(new ReadOperation() {
+ @Override
+ public CheckedFuture invoke(TransactionContext transactionContext) {
+ return transactionContext.dataExists(path);
+ }
+ });
}
+
private void checkModificationState() {
Preconditions.checkState(transactionType != TransactionType.READ_ONLY,
"Modification operation on read-only transaction is not allowed");
@@ -337,19 +300,12 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
LOG.debug("Tx {} write {}", identifier, path);
TransactionFutureCallback txFutureCallback = getOrCreateTxFutureCallback(path);
- TransactionContext transactionContext = txFutureCallback.getTransactionContext();
- if(transactionContext != null) {
- transactionContext.writeData(path, data);
- } else {
- // The shard Tx hasn't been created yet so add the Tx operation to the Tx Future
- // callback to be executed after the Tx is created.
- txFutureCallback.addTxOperationOnComplete(new TransactionOperation() {
- @Override
- public void invoke(TransactionContext transactionContext) {
- transactionContext.writeData(path, data);
- }
- });
- }
+ txFutureCallback.enqueueModifyOperation(new TransactionOperation() {
+ @Override
+ public void invoke(TransactionContext transactionContext) {
+ transactionContext.writeData(path, data);
+ }
+ });
}
@Override
@@ -360,19 +316,12 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
LOG.debug("Tx {} merge {}", identifier, path);
TransactionFutureCallback txFutureCallback = getOrCreateTxFutureCallback(path);
- TransactionContext transactionContext = txFutureCallback.getTransactionContext();
- if(transactionContext != null) {
- transactionContext.mergeData(path, data);
- } else {
- // The shard Tx hasn't been created yet so add the Tx operation to the Tx Future
- // callback to be executed after the Tx is created.
- txFutureCallback.addTxOperationOnComplete(new TransactionOperation() {
- @Override
- public void invoke(TransactionContext transactionContext) {
- transactionContext.mergeData(path, data);
- }
- });
- }
+ txFutureCallback.enqueueModifyOperation(new TransactionOperation() {
+ @Override
+ public void invoke(TransactionContext transactionContext) {
+ transactionContext.mergeData(path, data);
+ }
+ });
}
@Override
@@ -383,19 +332,12 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
LOG.debug("Tx {} delete {}", identifier, path);
TransactionFutureCallback txFutureCallback = getOrCreateTxFutureCallback(path);
- TransactionContext transactionContext = txFutureCallback.getTransactionContext();
- if(transactionContext != null) {
- transactionContext.deleteData(path);
- } else {
- // The shard Tx hasn't been created yet so add the Tx operation to the Tx Future
- // callback to be executed after the Tx is created.
- txFutureCallback.addTxOperationOnComplete(new TransactionOperation() {
- @Override
- public void invoke(TransactionContext transactionContext) {
- transactionContext.deleteData(path);
- }
- });
- }
+ txFutureCallback.enqueueModifyOperation(new TransactionOperation() {
+ @Override
+ public void invoke(TransactionContext transactionContext) {
+ transactionContext.deleteData(path);
+ }
+ });
}
@Override
@@ -412,35 +354,45 @@ public class TransactionProxy implements DOMStoreReadWriteTransaction {
for(TransactionFutureCallback txFutureCallback : txFutureCallbackMap.values()) {
- LOG.debug("Tx {} Readying transaction for shard {}", identifier,
- txFutureCallback.getShardName());
+ LOG.debug("Tx {} Readying transaction for shard {} chain {}", identifier,
+ txFutureCallback.getShardName(), transactionChainId);
- TransactionContext transactionContext = txFutureCallback.getTransactionContext();
- if(transactionContext != null) {
- cohortFutures.add(transactionContext.readyTransaction());
- } else {
- // The shard Tx hasn't been created yet so create a promise to ready the Tx later
- // after it's created.
- final Promise cohortPromise = akka.dispatch.Futures.promise();
- txFutureCallback.addTxOperationOnComplete(new TransactionOperation() {
- @Override
- public void invoke(TransactionContext transactionContext) {
- cohortPromise.completeWith(transactionContext.readyTransaction());
- }
- });
+ Future future = txFutureCallback.enqueueFutureOperation(new FutureOperation() {
+ @Override
+ public Future invoke(TransactionContext transactionContext) {
+ return transactionContext.readyTransaction();
+ }
+ });
- cohortFutures.add(cohortPromise.future());
- }
+ cohortFutures.add(future);
}
- if(transactionChainProxy != null){
- transactionChainProxy.onTransactionReady(cohortFutures);
- }
+ onTransactionReady(cohortFutures);
return new ThreePhaseCommitCohortProxy(actorContext, cohortFutures,
identifier.toString());
}
+ /**
+ * Method for derived classes to be notified when the transaction has been readied.
+ *
+ * @param cohortFutures the cohort Futures for each shard transaction.
+ */
+ 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