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.Collections;
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.controller.cluster.access.concepts.MemberName;
import org.opendaylight.controller.cluster.common.actor.AbstractUntypedPersistentActorWithMetering;
import org.opendaylight.controller.cluster.common.actor.Dispatchers;
-import org.opendaylight.controller.cluster.datastore.AbstractDataStore;
import org.opendaylight.controller.cluster.datastore.ClusterWrapper;
import org.opendaylight.controller.cluster.datastore.DatastoreContext;
import org.opendaylight.controller.cluster.datastore.DatastoreContext.Builder;
import org.opendaylight.controller.cluster.raft.messages.ServerChangeStatus;
import org.opendaylight.controller.cluster.raft.messages.ServerRemoved;
import org.opendaylight.controller.cluster.raft.policy.DisableElectionsRaftPolicy;
-import org.opendaylight.controller.cluster.sharding.PrefixedShardConfigUpdateHandler;
-import org.opendaylight.controller.cluster.sharding.messages.InitConfigListener;
-import org.opendaylight.controller.cluster.sharding.messages.PrefixShardCreated;
-import org.opendaylight.controller.cluster.sharding.messages.PrefixShardRemoved;
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;
private final Set<Consumer<String>> shardAvailabilityCallbacks = new HashSet<>();
private final String persistenceId;
- private final AbstractDataStore dataStore;
-
- private PrefixedShardConfigUpdateHandler configUpdateHandler;
ShardManager(final AbstractShardManagerCreator<?> builder) {
this.cluster = builder.getCluster();
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();
"shard-manager-" + this.type,
datastoreContextFactory.getBaseDatastoreContext().getDataStoreMXBeanType());
shardManagerMBean.registerMBean();
-
- dataStore = builder.getDistributedDataStore();
}
@Override
onAddShardReplica((AddShardReplica) message);
} else if (message instanceof AddPrefixShardReplica) {
onAddPrefixShardReplica((AddPrefixShardReplica) message);
- } else if (message instanceof PrefixShardCreated) {
- onPrefixShardCreated((PrefixShardCreated) message);
- } else if (message instanceof PrefixShardRemoved) {
- onPrefixShardRemoved((PrefixShardRemoved) message);
- } else if (message instanceof InitConfigListener) {
- onInitConfigListener();
} else if (message instanceof ForwardedAddServerReply) {
ForwardedAddServerReply msg = (ForwardedAddServerReply)message;
onAddServerReply(msg.shardInfo, msg.addServerReply, getSender(), msg.leaderPath,
} 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) {
getSender().tell(new GetShardRoleReply(shardInformation.getRole()), ActorRef.noSender());
}
- private void onInitConfigListener() {
- LOG.debug("{}: Initializing config listener on {}", persistenceId(), cluster.getCurrentMemberName());
-
- final org.opendaylight.mdsal.common.api.LogicalDatastoreType datastoreType =
- org.opendaylight.mdsal.common.api.LogicalDatastoreType
- .valueOf(datastoreContextFactory.getBaseDatastoreContext().getLogicalStoreType().name());
-
- if (configUpdateHandler != null) {
- configUpdateHandler.close();
- }
-
- configUpdateHandler = new PrefixedShardConfigUpdateHandler(self(), cluster.getCurrentMemberName());
- 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) {
}
}
+ @SuppressFBWarnings(value = "UPM_UNCALLED_PRIVATE_METHOD",
+ justification = "https://github.com/spotbugs/spotbugs/issues/811")
private void removePrefixShardReplica(final RemovePrefixShardReplica contextMessage, final String shardName,
final String primaryPath, final ActorRef sender) {
if (isShardReplicaOperationInProgress(shardName, sender)) {
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) {
}, new Dispatchers(context().system().dispatchers()).getDispatcher(Dispatchers.DispatcherType.Client));
}
+ @SuppressFBWarnings(value = "UPM_UNCALLED_PRIVATE_METHOD",
+ justification = "https://github.com/spotbugs/spotbugs/issues/811")
private void removeShardReplica(final RemoveShardReplica contextMessage, final String shardName,
final String primaryPath, final ActorRef sender) {
if (isShardReplicaOperationInProgress(shardName, sender)) {
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 onPrefixShardCreated(final PrefixShardCreated message) {
- LOG.debug("{}: onPrefixShardCreated: {}", persistenceId(), message);
-
- final PrefixShardConfiguration config = message.getConfiguration();
- final ShardIdentifier shardId = getShardIdentifier(cluster.getCurrentMemberName(),
- ClusterUtils.getCleanShardName(config.getPrefix().getRootIdentifier()));
- final String shardName = shardId.getShardName();
-
- if (isPreviousShardActorStopInProgress(shardName, message)) {
- return;
- }
-
- if (localShards.containsKey(shardName)) {
- LOG.debug("{}: Received create for an already existing shard {}", persistenceId(), shardName);
- final PrefixShardConfiguration existing =
- configuration.getAllPrefixShardConfigurations().get(config.getPrefix());
-
- if (existing != null && existing.equals(config)) {
- // we don't have to do nothing here
- return;
- }
- }
-
- doCreatePrefixShard(config, shardId, shardName);
- }
-
+ @SuppressFBWarnings(value = "UPM_UNCALLED_PRIVATE_METHOD",
+ justification = "https://github.com/spotbugs/spotbugs/issues/811")
private boolean isPreviousShardActorStopInProgress(final String shardName, final Object messageToDefer) {
final CompositeOnComplete<Boolean> stopOnComplete = shardActorsStopping.get(shardName);
if (stopOnComplete == null) {
return true;
}
- private void doCreatePrefixShard(final PrefixShardConfiguration config, final ShardIdentifier shardId,
- final String shardName) {
- configuration.addPrefixShardConfiguration(config);
-
- final Builder builder = newShardDatastoreContextBuilder(shardName);
- builder.logicalStoreType(config.getPrefix().getDatastoreType())
- .storeRoot(config.getPrefix().getRootIdentifier());
- DatastoreContext shardDatastoreContext = builder.build();
-
- final Map<String, String> peerAddresses = getPeerAddresses(shardName);
- final boolean isActiveMember = true;
-
- LOG.debug("{} doCreatePrefixShard: shardId: {}, memberNames: {}, peerAddresses: {}, isActiveMember: {}",
- persistenceId(), shardId, config.getShardMemberNames(), peerAddresses, isActiveMember);
-
- final ShardInformation info = new ShardInformation(shardName, shardId, peerAddresses,
- shardDatastoreContext, Shard.builder(), peerAddressResolver);
- info.setActiveMember(isActiveMember);
- localShards.put(info.getShardName(), info);
-
- if (schemaContext != null) {
- info.setSchemaContext(schemaContext);
- info.setActor(newShardActor(info));
- }
- }
-
- private void onPrefixShardRemoved(final PrefixShardRemoved message) {
- LOG.debug("{}: onPrefixShardRemoved : {}", persistenceId(), message);
-
- final DOMDataTreeIdentifier prefix = message.getPrefix();
- final ShardIdentifier shardId = getShardIdentifier(cluster.getCurrentMemberName(),
- ClusterUtils.getCleanShardName(prefix.getRootIdentifier()));
-
- configuration.removePrefixShardConfiguration(prefix);
- removeShard(shardId);
- }
-
private void doCreateShard(final CreateShard createShard) {
final ModuleShardConfiguration moduleShardConfig = createShard.getModuleShardConfig();
final String shardName = moduleShardConfig.getShardName();
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) {
}
@VisibleForTesting
- ShardInformation createShardInfoFor(String shardName, ShardIdentifier shardId,
- Map<String, String> peerAddresses,
- DatastoreContext datastoreContext,
- Map<String, DatastoreSnapshot.ShardSnapshot> shardSnapshots) {
+ ShardInformation createShardInfoFor(final String shardName, final ShardIdentifier shardId,
+ final Map<String, String> peerAddresses,
+ final DatastoreContext datastoreContext,
+ final Map<String, DatastoreSnapshot.ShardSnapshot> shardSnapshots) {
return new ShardInformation(shardName, shardId, peerAddresses,
datastoreContext, Shard.builder().restoreFromSnapshot(shardSnapshots.get(shardName)),
peerAddressResolver);
});
}
+ @SuppressFBWarnings(value = "UPM_UNCALLED_PRIVATE_METHOD",
+ justification = "https://github.com/spotbugs/spotbugs/issues/811")
private void sendLocalReplicaAlreadyExistsReply(final String shardName, final ActorRef sender) {
LOG.debug("{}: Local shard {} already exists", persistenceId(), shardName);
sender.tell(new Status.Failure(new AlreadyExistsException(
String.format("Local shard %s already exists", shardName))), getSelf());
}
+ @SuppressFBWarnings(value = "UPM_UNCALLED_PRIVATE_METHOD",
+ justification = "https://github.com/spotbugs/spotbugs/issues/811")
private void addPrefixShard(final String shardName, final YangInstanceIdentifier shardPrefix,
final RemotePrimaryShardFound response, final ActorRef sender) {
if (isShardReplicaOperationInProgress(shardName, sender)) {
execAddShard(shardName, shardInfo, response, removeShardOnFailure, sender);
}
+ @SuppressFBWarnings(value = "UPM_UNCALLED_PRIVATE_METHOD",
+ justification = "https://github.com/spotbugs/spotbugs/issues/811")
private void addShard(final String shardName, final RemotePrimaryShardFound response, final ActorRef sender) {
if (isShardReplicaOperationInProgress(shardName, sender)) {
return;
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);