package org.opendaylight.controller.cluster.datastore;
import static akka.pattern.Patterns.ask;
-import akka.actor.ActorPath;
import akka.actor.ActorRef;
import akka.actor.Address;
import akka.actor.Cancellable;
}
private void onShutDown() {
- Shutdown shutdown = new Shutdown();
List<Future<Boolean>> stopFutures = new ArrayList<>(localShards.size());
for (ShardInformation info : localShards.values()) {
if (info.getActor() != null) {
LOG.debug("{}: Issuing gracefulStop to shard {}", persistenceId(), info.getShardId());
FiniteDuration duration = info.getDatastoreContext().getShardRaftConfig().getElectionTimeOutInterval().$times(2);
- stopFutures.add(Patterns.gracefulStop(info.getActor(), duration, shutdown));
+ stopFutures.add(Patterns.gracefulStop(info.getActor(), duration, Shutdown.INSTANCE));
}
}
return;
} else if(shardInformation.getActor() != null) {
LOG.debug("{} : Sending Shutdown to Shard actor {}", persistenceId(), shardInformation.getActor());
- shardInformation.getActor().tell(new Shutdown(), self());
+ shardInformation.getActor().tell(Shutdown.INSTANCE, self());
}
LOG.debug("{} : Local Shard replica for shard {} has been removed", persistenceId(), shardId.getShardName());
persistShardList();
}
private void memberRemoved(ClusterEvent.MemberRemoved message) {
- String memberName = message.member().roles().head();
+ String memberName = message.member().roles().iterator().next();
LOG.debug("{}: Received MemberRemoved: memberName: {}, address: {}", persistenceId(), memberName,
message.member().address());
}
private void memberExited(ClusterEvent.MemberExited message) {
- String memberName = message.member().roles().head();
+ String memberName = message.member().roles().iterator().next();
LOG.debug("{}: Received MemberExited: memberName: {}, address: {}", persistenceId(), memberName,
message.member().address());
}
private void memberUp(ClusterEvent.MemberUp message) {
- String memberName = message.member().roles().head();
+ String memberName = message.member().roles().iterator().next();
LOG.debug("{}: Received MemberUp: memberName: {}, address: {}", persistenceId(), memberName,
message.member().address());
}
private void memberReachable(ClusterEvent.ReachableMember message) {
- String memberName = message.member().roles().head();
+ String memberName = message.member().roles().iterator().next();
LOG.debug("Received ReachableMember: memberName {}, address: {}", memberName, message.member().address());
addPeerAddress(memberName, message.member().address());
}
private void memberUnreachable(ClusterEvent.UnreachableMember message) {
- String memberName = message.member().roles().head();
+ String memberName = message.member().roles().iterator().next();
LOG.debug("Received UnreachableMember: memberName {}, address: {}", memberName, message.member().address());
markMemberUnavailable(memberName);
continue;
}
- LOG.debug("{}: findPrimary for {} forwarding to remote ShardManager {}", persistenceId(),
- shardName, address);
+ LOG.debug("{}: findPrimary for {} forwarding to remote ShardManager {}, visitedAddresses: {}",
+ persistenceId(), shardName, address, visitedAddresses);
getContext().actorSelection(address).forward(new RemoteFindPrimary(shardName,
message.isWaitUntilReady(), visitedAddresses), getContext());
}
}
- private Exception getServerChangeException(Class<?> serverChange, ServerChangeStatus serverChangeStatus,
+ private static Exception getServerChangeException(Class<?> serverChange, ServerChangeStatus serverChangeStatus,
String leaderPath, ShardIdentifier shardId) {
Exception failure;
switch (serverChangeStatus) {
private void onSaveSnapshotSuccess (SaveSnapshotSuccess successMessage) {
LOG.debug ("{} saved ShardManager snapshot successfully. Deleting the prev snapshot if available",
persistenceId());
- deleteSnapshots(new SnapshotSelectionCriteria(scala.Long.MaxValue(), (successMessage.metadata().timestamp() - 1)));
+ deleteSnapshots(new SnapshotSelectionCriteria(scala.Long.MaxValue(), successMessage.metadata().timestamp() - 1,
+ 0, 0));
}
private static class ForwardedAddServerReply {
private final ShardIdentifier shardId;
private final String shardName;
private ActorRef actor;
- private ActorPath actorPath;
private final Map<String, String> initialPeerAddresses;
private Optional<DataTree> localShardDataTree;
private boolean leaderAvailable = false;
return actor;
}
- ActorPath getActorPath() {
- return actorPath;
- }
-
void setActor(ActorRef actor) {
this.actor = actor;
- this.actorPath = actor.path();
}
ShardIdentifier getShardId() {
private CountDownLatch waitTillReadyCountdownLatch;
private PrimaryShardInfoFutureCache primaryShardInfoCache;
private DatastoreSnapshot restoreFromSnapshot;
-
private volatile boolean sealed;
@SuppressWarnings("unchecked")