import akka.persistence.SnapshotSelectionCriteria;
import akka.util.Timeout;
import com.google.common.annotations.VisibleForTesting;
+import com.google.common.util.concurrent.SettableFuture;
import edu.umd.cs.findbugs.annotations.SuppressFBWarnings;
import java.util.ArrayList;
import java.util.Collection;
import java.util.Map;
import java.util.Map.Entry;
import java.util.Set;
-import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
import java.util.function.Consumer;
import org.opendaylight.mdsal.dom.api.DOMDataTreeIdentifier;
import org.opendaylight.yangtools.concepts.Registration;
import org.opendaylight.yangtools.yang.data.api.YangInstanceIdentifier;
-import org.opendaylight.yangtools.yang.model.api.SchemaContext;
+import org.opendaylight.yangtools.yang.model.api.EffectiveModelContext;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import scala.concurrent.ExecutionContext;
private DatastoreContextFactory datastoreContextFactory;
- private final CountDownLatch waitTillReadyCountdownLatch;
+ private final SettableFuture<Void> readinessFuture;
private final PrimaryShardInfoFutureCache primaryShardInfoCache;
@VisibleForTesting
final ShardPeerAddressResolver peerAddressResolver;
- private SchemaContext schemaContext;
+ private EffectiveModelContext schemaContext;
private DatastoreSnapshot restoreFromSnapshot;
this.type = datastoreContextFactory.getBaseDatastoreContext().getDataStoreName();
this.shardDispatcherPath =
new Dispatchers(context().system().dispatchers()).getDispatcherPath(Dispatchers.DispatcherType.Shard);
- this.waitTillReadyCountdownLatch = builder.getWaitTillReadyCountDownLatch();
+ this.readinessFuture = builder.getReadinessFuture();
this.primaryShardInfoCache = builder.getPrimaryShardInfoCache();
this.restoreFromSnapshot = builder.getRestoreFromSnapshot();
} else if (message instanceof WrappedShardResponse) {
onWrappedShardResponse((WrappedShardResponse) message);
} else if (message instanceof GetSnapshot) {
- onGetSnapshot();
+ onGetSnapshot((GetSnapshot) message);
} else if (message instanceof ServerRemoved) {
onShardReplicaRemoved((ServerRemoved) message);
} else if (message instanceof ChangeShardMembersVotingStatus) {
configUpdateHandler.initListener(dataStore, datastoreType);
}
- private void onShutDown() {
+ void onShutDown() {
List<Future<Boolean>> stopFutures = new ArrayList<>(localShards.size());
for (ShardInformation info : localShards.values()) {
if (info.getActor() != null) {
Future<Object> futureObj = ask(getContext().actorSelection(primaryPath),
new RemoveServer(shardId.toString()), removeServerTimeout);
- futureObj.onComplete(new OnComplete<Object>() {
+ futureObj.onComplete(new OnComplete<>() {
@Override
public void onComplete(final Throwable failure, final Object response) {
if (failure != null) {
Future<Object> futureObj = ask(getContext().actorSelection(primaryPath),
new RemoveServer(shardId.toString()), removeServerTimeout);
- futureObj.onComplete(new OnComplete<Object>() {
+ futureObj.onComplete(new OnComplete<>() {
@Override
public void onComplete(final Throwable failure, final Object response) {
if (failure != null) {
final Future<Boolean> stopFuture = Patterns.gracefulStop(shardActor,
FiniteDuration.apply(timeoutInMS, TimeUnit.MILLISECONDS), Shutdown.INSTANCE);
- final CompositeOnComplete<Boolean> onComplete = new CompositeOnComplete<Boolean>() {
+ final CompositeOnComplete<Boolean> onComplete = new CompositeOnComplete<>() {
@Override
public void onComplete(final Throwable failure, final Boolean result) {
if (failure == null) {
persistShardList();
}
- private void onGetSnapshot() {
+ private void onGetSnapshot(final GetSnapshot getSnapshot) {
LOG.debug("{}: onGetSnapshot", persistenceId());
List<String> notInitialized = null;
datastoreContextFactory.getBaseDatastoreContext().getShardInitializationTimeout().duration()));
for (ShardInformation shardInfo: localShards.values()) {
- shardInfo.getActor().tell(GetSnapshot.INSTANCE, replyActor);
+ shardInfo.getActor().tell(getSnapshot, replyActor);
}
}
private void checkReady() {
if (isReadyWithLeaderId()) {
- LOG.info("{}: All Shards are ready - data store {} is ready, available count is {}",
- persistenceId(), type, waitTillReadyCountdownLatch.getCount());
-
- waitTillReadyCountdownLatch.countDown();
+ LOG.info("{}: All Shards are ready - data store {} is ready", persistenceId(), type);
+ readinessFuture.set(null);
}
}
* @param message the message to send
*/
private void updateSchemaContext(final Object message) {
- schemaContext = ((UpdateSchemaContext) message).getSchemaContext();
+ schemaContext = ((UpdateSchemaContext) message).getEffectiveModelContext();
LOG.debug("Got updated SchemaContext: # of modules {}", schemaContext.getModules().size());
.getShardInitializationTimeout().duration().$times(2));
Future<Object> futureObj = ask(getSelf(), new FindPrimary(shardName, true), findPrimaryTimeout);
- futureObj.onComplete(new OnComplete<Object>() {
+ futureObj.onComplete(new OnComplete<>() {
@Override
public void onComplete(final Throwable failure, final Object response) {
if (failure != null) {
final Future<Object> futureObj = ask(getContext().actorSelection(response.getPrimaryPath()),
new AddServer(shardInfo.getShardId().toString(), localShardAddress, true), addServerTimeout);
- futureObj.onComplete(new OnComplete<Object>() {
+ futureObj.onComplete(new OnComplete<>() {
@Override
public void onComplete(final Throwable failure, final Object addServerResponse) {
if (failure != null) {
Future<Object> future = ask(localShardFound.getPath(), GetOnDemandRaftState.INSTANCE,
Timeout.apply(30, TimeUnit.SECONDS));
- future.onComplete(new OnComplete<Object>() {
+ future.onComplete(new OnComplete<>() {
@Override
public void onComplete(final Throwable failure, final Object response) {
if (failure != null) {
.getShardInitializationTimeout().duration().$times(2));
Future<Object> futureObj = ask(getSelf(), new FindLocalShard(shardName, true), findLocalTimeout);
- futureObj.onComplete(new OnComplete<Object>() {
+ futureObj.onComplete(new OnComplete<>() {
@Override
public void onComplete(final Throwable failure, final Object response) {
if (failure != null) {
Timeout timeout = new Timeout(datastoreContext.getShardLeaderElectionTimeout().duration().$times(2));
Future<Object> futureObj = ask(shardActorRef, changeServersVotingStatus, timeout);
- futureObj.onComplete(new OnComplete<Object>() {
+ futureObj.onComplete(new OnComplete<>() {
@Override
public void onComplete(final Throwable failure, final Object response) {
shardReplicaOperationsInProgress.remove(shardName);