4644a2d79807bf36d069fa010ce561e0232cbcf9
[controller.git] / opendaylight / md-sal / sal-akka-raft / src / main / java / org / opendaylight / controller / cluster / raft / RaftActorContextImpl.java
1 /*
2  * Copyright (c) 2014 Cisco Systems, Inc. and others.  All rights reserved.
3  *
4  * This program and the accompanying materials are made available under the
5  * terms of the Eclipse Public License v1.0 which accompanies this distribution,
6  * and is available at http://www.eclipse.org/legal/epl-v10.html
7  */
8
9 package org.opendaylight.controller.cluster.raft;
10
11 import akka.actor.ActorContext;
12 import akka.actor.ActorRef;
13 import akka.actor.ActorSelection;
14 import akka.actor.ActorSystem;
15 import akka.actor.Props;
16 import akka.cluster.Cluster;
17 import com.google.common.annotations.VisibleForTesting;
18 import com.google.common.base.Preconditions;
19 import java.util.ArrayList;
20 import java.util.Collection;
21 import java.util.HashMap;
22 import java.util.HashSet;
23 import java.util.List;
24 import java.util.Map;
25 import java.util.Optional;
26 import java.util.Set;
27 import java.util.function.Consumer;
28 import java.util.function.LongSupplier;
29 import javax.annotation.Nonnull;
30 import javax.annotation.Nullable;
31 import org.opendaylight.controller.cluster.DataPersistenceProvider;
32 import org.opendaylight.controller.cluster.io.FileBackedOutputStreamFactory;
33 import org.opendaylight.controller.cluster.raft.base.messages.ApplyState;
34 import org.opendaylight.controller.cluster.raft.behaviors.RaftActorBehavior;
35 import org.opendaylight.controller.cluster.raft.persisted.ServerConfigurationPayload;
36 import org.opendaylight.controller.cluster.raft.persisted.ServerInfo;
37 import org.opendaylight.controller.cluster.raft.policy.RaftPolicy;
38 import org.slf4j.Logger;
39
40 /**
41  * Implementation of the RaftActorContext interface.
42  *
43  * @author Moiz Raja
44  * @author Thomas Pantelis
45  */
46 public class RaftActorContextImpl implements RaftActorContext {
47     private static final LongSupplier JVM_MEMORY_RETRIEVER = () -> Runtime.getRuntime().maxMemory();
48
49     private final ActorRef actor;
50
51     private final ActorContext context;
52
53     private final String id;
54
55     private final ElectionTerm termInformation;
56
57     private long commitIndex;
58
59     private long lastApplied;
60
61     private ReplicatedLog replicatedLog;
62
63     private final Map<String, PeerInfo> peerInfoMap = new HashMap<>();
64
65     private final Logger log;
66
67     private ConfigParams configParams;
68
69     private boolean dynamicServerConfiguration = false;
70
71     @VisibleForTesting
72     private LongSupplier totalMemoryRetriever = JVM_MEMORY_RETRIEVER;
73
74     // Snapshot manager will need to be created on demand as it needs raft actor context which cannot
75     // be passed to it in the constructor
76     private SnapshotManager snapshotManager;
77
78     private final DataPersistenceProvider persistenceProvider;
79
80     private short payloadVersion;
81
82     private boolean votingMember = true;
83
84     private RaftActorBehavior currentBehavior;
85
86     private int numVotingPeers = -1;
87
88     private Optional<Cluster> cluster;
89
90     private final Consumer<ApplyState> applyStateConsumer;
91
92     private final FileBackedOutputStreamFactory fileBackedOutputStreamFactory;
93
94     private RaftActorLeadershipTransferCohort leadershipTransferCohort;
95
96     public RaftActorContextImpl(ActorRef actor, ActorContext context, String id,
97             @Nonnull ElectionTerm termInformation, long commitIndex, long lastApplied,
98             @Nonnull Map<String, String> peerAddresses,
99             @Nonnull ConfigParams configParams, @Nonnull DataPersistenceProvider persistenceProvider,
100             @Nonnull Consumer<ApplyState> applyStateConsumer, @Nonnull Logger logger) {
101         this.actor = actor;
102         this.context = context;
103         this.id = id;
104         this.termInformation = Preconditions.checkNotNull(termInformation);
105         this.commitIndex = commitIndex;
106         this.lastApplied = lastApplied;
107         this.configParams = Preconditions.checkNotNull(configParams);
108         this.persistenceProvider = Preconditions.checkNotNull(persistenceProvider);
109         this.log = Preconditions.checkNotNull(logger);
110         this.applyStateConsumer = Preconditions.checkNotNull(applyStateConsumer);
111
112         fileBackedOutputStreamFactory = new FileBackedOutputStreamFactory(
113                 configParams.getFileBackedStreamingThreshold(), configParams.getTempFileDirectory());
114
115         for (Map.Entry<String, String> e: Preconditions.checkNotNull(peerAddresses).entrySet()) {
116             peerInfoMap.put(e.getKey(), new PeerInfo(e.getKey(), e.getValue(), VotingState.VOTING));
117         }
118     }
119
120     @VisibleForTesting
121     public void setPayloadVersion(short payloadVersion) {
122         this.payloadVersion = payloadVersion;
123     }
124
125     @Override
126     public short getPayloadVersion() {
127         return payloadVersion;
128     }
129
130     public void setConfigParams(ConfigParams configParams) {
131         this.configParams = configParams;
132     }
133
134     @Override
135     public ActorRef actorOf(Props props) {
136         return context.actorOf(props);
137     }
138
139     @Override
140     public ActorSelection actorSelection(String path) {
141         return context.actorSelection(path);
142     }
143
144     @Override
145     public String getId() {
146         return id;
147     }
148
149     @Override
150     public ActorRef getActor() {
151         return actor;
152     }
153
154     @Override
155     @SuppressWarnings("checkstyle:IllegalCatch")
156     public Optional<Cluster> getCluster() {
157         if (cluster == null) {
158             try {
159                 cluster = Optional.of(Cluster.get(getActorSystem()));
160             } catch (Exception e) {
161                 // An exception means there's no cluster configured. This will only happen in unit tests.
162                 log.debug("{}: Could not obtain Cluster: {}", getId(), e);
163                 cluster = Optional.empty();
164             }
165         }
166
167         return cluster;
168     }
169
170     @Override
171     public ElectionTerm getTermInformation() {
172         return termInformation;
173     }
174
175     @Override
176     public long getCommitIndex() {
177         return commitIndex;
178     }
179
180     @Override public void setCommitIndex(long commitIndex) {
181         this.commitIndex = commitIndex;
182     }
183
184     @Override
185     public long getLastApplied() {
186         return lastApplied;
187     }
188
189     @Override
190     public void setLastApplied(long lastApplied) {
191         log.debug("{}: Moving last applied index from {} to {}", id, this.lastApplied, lastApplied);
192         this.lastApplied = lastApplied;
193     }
194
195     @Override
196     public void setReplicatedLog(ReplicatedLog replicatedLog) {
197         this.replicatedLog = replicatedLog;
198     }
199
200     @Override
201     public ReplicatedLog getReplicatedLog() {
202         return replicatedLog;
203     }
204
205     @Override public ActorSystem getActorSystem() {
206         return context.system();
207     }
208
209     @Override public Logger getLogger() {
210         return this.log;
211     }
212
213     @Override
214     public Collection<String> getPeerIds() {
215         return peerInfoMap.keySet();
216     }
217
218     @Override
219     public Collection<PeerInfo> getPeers() {
220         return peerInfoMap.values();
221     }
222
223     @Override
224     public PeerInfo getPeerInfo(String peerId) {
225         return peerInfoMap.get(peerId);
226     }
227
228     @Override
229     public String getPeerAddress(String peerId) {
230         String peerAddress;
231         PeerInfo peerInfo = peerInfoMap.get(peerId);
232         if (peerInfo != null) {
233             peerAddress = peerInfo.getAddress();
234             if (peerAddress == null) {
235                 peerAddress = configParams.getPeerAddressResolver().resolve(peerId);
236                 peerInfo.setAddress(peerAddress);
237             }
238         } else {
239             peerAddress = configParams.getPeerAddressResolver().resolve(peerId);
240         }
241
242         return peerAddress;
243     }
244
245     @Override
246     public void updatePeerIds(ServerConfigurationPayload serverConfig) {
247         votingMember = true;
248         boolean foundSelf = false;
249         Set<String> currentPeers = new HashSet<>(this.getPeerIds());
250         for (ServerInfo server : serverConfig.getServerConfig()) {
251             if (getId().equals(server.getId())) {
252                 foundSelf = true;
253                 if (!server.isVoting()) {
254                     votingMember = false;
255                 }
256             } else {
257                 VotingState votingState = server.isVoting() ? VotingState.VOTING : VotingState.NON_VOTING;
258                 if (!currentPeers.contains(server.getId())) {
259                     this.addToPeers(server.getId(), null, votingState);
260                 } else {
261                     this.getPeerInfo(server.getId()).setVotingState(votingState);
262                     currentPeers.remove(server.getId());
263                 }
264             }
265         }
266
267         for (String peerIdToRemove : currentPeers) {
268             this.removePeer(peerIdToRemove);
269         }
270
271         if (!foundSelf) {
272             votingMember = false;
273         }
274
275         log.debug("{}: Updated server config: isVoting: {}, peers: {}", id, votingMember, peerInfoMap.values());
276
277         setDynamicServerConfigurationInUse();
278     }
279
280     @Override public ConfigParams getConfigParams() {
281         return configParams;
282     }
283
284     @Override
285     public void addToPeers(String peerId, String address, VotingState votingState) {
286         peerInfoMap.put(peerId, new PeerInfo(peerId, address, votingState));
287         numVotingPeers = -1;
288     }
289
290     @Override
291     public void removePeer(String name) {
292         if (getId().equals(name)) {
293             votingMember = false;
294         } else {
295             peerInfoMap.remove(name);
296             numVotingPeers = -1;
297         }
298     }
299
300     @Override public ActorSelection getPeerActorSelection(String peerId) {
301         String peerAddress = getPeerAddress(peerId);
302         if (peerAddress != null) {
303             return actorSelection(peerAddress);
304         }
305         return null;
306     }
307
308     @Override
309     public void setPeerAddress(String peerId, String peerAddress) {
310         PeerInfo peerInfo = peerInfoMap.get(peerId);
311         if (peerInfo != null) {
312             log.info("Peer address for peer {} set to {}", peerId, peerAddress);
313             peerInfo.setAddress(peerAddress);
314         }
315     }
316
317     @Override
318     public SnapshotManager getSnapshotManager() {
319         if (snapshotManager == null) {
320             snapshotManager = new SnapshotManager(this, log);
321         }
322         return snapshotManager;
323     }
324
325     @Override
326     public long getTotalMemory() {
327         return totalMemoryRetriever.getAsLong();
328     }
329
330     @Override
331     public void setTotalMemoryRetriever(LongSupplier retriever) {
332         totalMemoryRetriever = retriever == null ? JVM_MEMORY_RETRIEVER : retriever;
333     }
334
335     @Override
336     public boolean hasFollowers() {
337         return !getPeerIds().isEmpty();
338     }
339
340     @Override
341     public DataPersistenceProvider getPersistenceProvider() {
342         return persistenceProvider;
343     }
344
345
346     @Override
347     public RaftPolicy getRaftPolicy() {
348         return configParams.getRaftPolicy();
349     }
350
351     @Override
352     public boolean isDynamicServerConfigurationInUse() {
353         return dynamicServerConfiguration;
354     }
355
356     @Override
357     public void setDynamicServerConfigurationInUse() {
358         this.dynamicServerConfiguration = true;
359     }
360
361     @Override
362     public ServerConfigurationPayload getPeerServerInfo(boolean includeSelf) {
363         if (!isDynamicServerConfigurationInUse()) {
364             return null;
365         }
366         Collection<PeerInfo> peers = getPeers();
367         List<ServerInfo> newConfig = new ArrayList<>(peers.size() + 1);
368         for (PeerInfo peer: peers) {
369             newConfig.add(new ServerInfo(peer.getId(), peer.isVoting()));
370         }
371
372         if (includeSelf) {
373             newConfig.add(new ServerInfo(getId(), votingMember));
374         }
375
376         return new ServerConfigurationPayload(newConfig);
377     }
378
379     @Override
380     public boolean isVotingMember() {
381         return votingMember;
382     }
383
384     @Override
385     public boolean anyVotingPeers() {
386         if (numVotingPeers < 0) {
387             numVotingPeers = 0;
388             for (PeerInfo info: getPeers()) {
389                 if (info.isVoting()) {
390                     numVotingPeers++;
391                 }
392             }
393         }
394
395         return numVotingPeers > 0;
396     }
397
398     @Override
399     public RaftActorBehavior getCurrentBehavior() {
400         return currentBehavior;
401     }
402
403     void setCurrentBehavior(final RaftActorBehavior behavior) {
404         this.currentBehavior = Preconditions.checkNotNull(behavior);
405     }
406
407     @Override
408     public Consumer<ApplyState> getApplyStateConsumer() {
409         return applyStateConsumer;
410     }
411
412     @Override
413     public FileBackedOutputStreamFactory getFileBackedOutputStreamFactory() {
414         return fileBackedOutputStreamFactory;
415     }
416
417     @SuppressWarnings("checkstyle:IllegalCatch")
418     void close() {
419         if (currentBehavior != null) {
420             try {
421                 currentBehavior.close();
422             } catch (Exception e) {
423                 log.debug("{}: Error closing behavior {}", getId(), currentBehavior.state(), e);
424             }
425         }
426     }
427
428     @Override
429     @Nullable
430     public RaftActorLeadershipTransferCohort getRaftActorLeadershipTransferCohort() {
431         return leadershipTransferCohort;
432     }
433
434     @Override
435     public void setRaftActorLeadershipTransferCohort(
436             @Nullable RaftActorLeadershipTransferCohort leadershipTransferCohort) {
437         this.leadershipTransferCohort = leadershipTransferCohort;
438     }
439 }

©2013 OpenDaylight, A Linux Foundation Collaborative Project. All Rights Reserved.
OpenDaylight is a registered trademark of The OpenDaylight Project, Inc.
Linux Foundation and OpenDaylight are registered trademarks of the Linux Foundation.
Linux is a registered trademark of Linus Torvalds.