package org.opendaylight.controller.cluster.raft;
-import static com.google.common.base.Preconditions.checkState;
import akka.actor.ActorRef;
import akka.actor.ActorSelection;
import akka.actor.ActorSystem;
import akka.actor.Props;
import akka.actor.UntypedActorContext;
+
import com.google.common.annotations.VisibleForTesting;
import com.google.common.base.Supplier;
+import com.google.common.collect.Maps;
+
+import java.util.Collection;
+import java.util.HashMap;
import java.util.Map;
+
+import org.opendaylight.controller.cluster.DataPersistenceProvider;
+import org.opendaylight.controller.cluster.raft.policy.RaftPolicy;
import org.slf4j.Logger;
public class RaftActorContextImpl implements RaftActorContext {
// be passed to it in the constructor
private SnapshotManager snapshotManager;
+ private final DataPersistenceProvider persistenceProvider;
+
+ private short payloadVersion;
+
public RaftActorContextImpl(ActorRef actor, UntypedActorContext context, String id,
ElectionTerm termInformation, long commitIndex, long lastApplied, Map<String, String> peerAddresses,
- ConfigParams configParams, Logger logger) {
+ ConfigParams configParams, DataPersistenceProvider persistenceProvider, Logger logger) {
this.actor = actor;
this.context = context;
this.id = id;
this.termInformation = termInformation;
this.commitIndex = commitIndex;
this.lastApplied = lastApplied;
- this.peerAddresses = peerAddresses;
+ this.peerAddresses = Maps.newHashMap(peerAddresses);
this.configParams = configParams;
+ this.persistenceProvider = persistenceProvider;
this.LOG = logger;
}
+ void setPayloadVersion(short payloadVersion) {
+ this.payloadVersion = payloadVersion;
+ }
+
+ @Override
+ public short getPayloadVersion() {
+ return payloadVersion;
+ }
+
void setConfigParams(ConfigParams configParams) {
this.configParams = configParams;
}
return lastApplied;
}
- @Override public void setLastApplied(long lastApplied) {
+ @Override
+ public void setLastApplied(long lastApplied) {
this.lastApplied = lastApplied;
}
- @Override public void setReplicatedLog(ReplicatedLog replicatedLog) {
+ @Override
+ public void setReplicatedLog(ReplicatedLog replicatedLog) {
this.replicatedLog = replicatedLog;
}
- @Override public ReplicatedLog getReplicatedLog() {
+ @Override
+ public ReplicatedLog getReplicatedLog() {
return replicatedLog;
}
return this.LOG;
}
- @Override public Map<String, String> getPeerAddresses() {
- return peerAddresses;
+ @Override
+ public Map<String, String> getPeerAddresses() {
+ return new HashMap<String, String>(peerAddresses);
+ }
+
+ @Override
+ public Collection<String> getPeerIds() {
+ return peerAddresses.keySet();
}
@Override public String getPeerAddress(String peerId) {
- return peerAddresses.get(peerId);
+ String peerAddress = peerAddresses.get(peerId);
+ if(peerAddress == null) {
+ peerAddress = configParams.getPeerAddressResolver().resolve(peerId);
+ peerAddresses.put(peerId, peerAddress);
+ }
+
+ return peerAddress;
}
@Override public ConfigParams getConfigParams() {
return null;
}
- @Override public void setPeerAddress(String peerId, String peerAddress) {
- LOG.info("Peer address for peer {} set to {}", peerId, peerAddress);
- checkState(peerAddresses.containsKey(peerId), peerId + " is unknown");
-
- peerAddresses.put(peerId, peerAddress);
+ @Override
+ public void setPeerAddress(String peerId, String peerAddress) {
+ if(peerAddresses.containsKey(peerId)) {
+ LOG.info("Peer address for peer {} set to {}", peerId, peerAddress);
+ peerAddresses.put(peerId, peerAddress);
+ }
}
@Override
@Override
public boolean hasFollowers() {
- return getPeerAddresses().keySet().size() > 0;
+ return getPeerIds().size() > 0;
+ }
+
+ @Override
+ public DataPersistenceProvider getPersistenceProvider() {
+ return persistenceProvider;
+ }
+
+
+ @Override
+ public RaftPolicy getRaftPolicy() {
+ return configParams.getRaftPolicy();
}
}