From 48620394657dba9b2547b9047d881041223c3f29 Mon Sep 17 00:00:00 2001 From: coolvikas Date: Sat, 21 Mar 2020 17:38:50 +0530 Subject: [PATCH] some minor code changes fr code semantics --- src/raft/raft.go | 116 ++++++++++++++++++++++++++++------------------- 1 file changed, 69 insertions(+), 47 deletions(-) diff --git a/src/raft/raft.go b/src/raft/raft.go index 7d1d60c..f93da06 100644 --- a/src/raft/raft.go +++ b/src/raft/raft.go @@ -89,6 +89,7 @@ type Raft struct { MatchIndex []int CurrentState int Mtx sync.Mutex + // why do we need voteCount? VoteCount map[int]int ElecTimeOutChannel chan struct{} AppendEChannel chan struct{} @@ -221,17 +222,20 @@ func (rf *Raft) AppendEntries(args *AppendEntryArgs, reply *AppendEntryReply) { if rf.CurrentTerm < args.Term { rf.CurrentTerm = args.Term rep.Success = true - rep.Term = rf.CurrentTerm // not correct BTW + rep.Term = rf.CurrentTerm // not correct BTW => it is correct as per figure2 of paper, since we have already updated the term of the server we need to send updated term to leader rf.VotedFor = -1 rf.CurrentState = STATE_FOLLOWER } else if rf.CurrentTerm == args.Term { rep.Success = true + rep.Term = rf.CurrentTerm + // do we really need to change to FOLLOWER as this server might also be having its own election? rf.CurrentState = STATE_FOLLOWER } else { rep.Success = false rep.Term = rf.CurrentTerm } rf.Mtx.Unlock() + // incase we receive an heartbeat message from a leader in network we need to stop and restart the election Timer rf.ElecTimeOutChannel <- struct{}{} // just to stop the timer of Election } else { // TODO LATER for actual Append Entries RPC from the Leader @@ -267,17 +271,22 @@ func (rf *Raft) RequestVote(args *RequestVoteArgs, reply *RequestVoteReply) { rf.CurrentTerm = args.Term rf.VotedFor = -1 if rf.CurrentState == STATE_FOLLOWER || rf.CurrentState == STATE_CANDIDATE { - if rf.VotedFor == -1 || rf.VotedFor == args.CandidateID { + // if we have already voted for CandidateID in current term, then we cannot vote again because it amy be a duplicate RPC and leader might count twice if we reply + // if rf.VotedFor == -1 || rf.VotedFor == args.CandidateID { + if rf.VotedFor == -1 { + // if we do not have any log entry in the log then we can vote otherwise we need a check if len(rf.Log) == 0 { - //log.Println(rf.me, " voting for", args.CandidateID) + log.Println(rf.me, " voting for", args.CandidateID) Voted = true rf.VotedFor = args.CandidateID } else { - if rf.Log[len(rf.Log)-1].Term <= args.LastLogTerm { + if rf.Log[len(rf.Log)-1].Term < args.LastLogTerm { Voted = true + rf.VotedFor = args.CandidateID } else if rf.Log[len(rf.Log)-1].Term == args.LastLogTerm { if len(rf.Log) <= args.LastLogIndex { Voted = true + rf.VotedFor = args.CandidateID } } } @@ -285,19 +294,22 @@ func (rf *Raft) RequestVote(args *RequestVoteArgs, reply *RequestVoteReply) { } } } else if args.Term == rf.CurrentTerm { - // rf.VotedFor = -1 if rf.CurrentState == STATE_FOLLOWER { - if rf.VotedFor == -1 || rf.VotedFor == args.CandidateID { + // if rf.VotedFor == -1 || rf.VotedFor == args.CandidateID { + // if we have already voted for CandidateID in current term, then we cannot vote again because it amy be a duplicate RPC and leader might count twice if we reply + if rf.VotedFor == -1 { if len(rf.Log) == 0 { log.Println(rf.me, " voting for", args.CandidateID) Voted = true rf.VotedFor = args.CandidateID } else { - if rf.Log[len(rf.Log)-1].Term <= args.LastLogTerm { + if rf.Log[len(rf.Log)-1].Term < args.LastLogTerm { Voted = true + rf.VotedFor = args.CandidateID } else if rf.Log[len(rf.Log)-1].Term == args.LastLogTerm { if len(rf.Log) <= args.LastLogIndex { Voted = true + rf.VotedFor = args.CandidateID } } } @@ -309,9 +321,6 @@ func (rf *Raft) RequestVote(args *RequestVoteArgs, reply *RequestVoteReply) { } rep.VoteGranted = Voted rep.Term = rf.CurrentTerm - // if args.Term > rf.CurrentTerm { - // rf.CurrentTerm = args.Term - // } //log.Println(rf.me, " got RequestVote from", args.CandidateID, " and voted = ", Voted) //log.Println(rf.me, " :my currentState is = ", rf.CurrentState, " : and my currentTerm is = ", rf.CurrentTerm) *reply = rep @@ -420,8 +429,8 @@ func sendAppendEntries(rf *Raft, _nServer int, _currentTerm int, _prevLogIndex i for serverID := 0; serverID < _nServer; serverID++ { //log.Println("sent ", serverID) - var req = AppendEntryArgs{} - var rep = AppendEntryReply{} + req := AppendEntryArgs{} + rep := AppendEntryReply{} req.Term = _currentTerm req.LeaderID = me req.PrevLogIndex = _prevLogIndex @@ -464,6 +473,7 @@ func Make(peers []*labrpc.ClientEnd, me int, rf.VotedFor = -1 rand.Seed(time.Now().UTC().UnixNano()) + // why a go routine is required here? go func() { // kick off leader election periodically by sending out RequestVote RPCs // when it hasn't heard from another peer for a while @@ -477,33 +487,37 @@ func Make(peers []*labrpc.ClientEnd, me int, t1 := time.Now() for { select { + //timer expired case case <-d.Done(): // Election Timer expired. t2 := time.Now() //fmt.Println(me, " ]> Election Started after timer expired after ", t2.Sub(t1)) + // why we are getiing a new context here? d, cancel = context.WithTimeout(context.Background(), time.Millisecond*time.Duration(rand.Intn(MAX_TIMEOUT-MIN_TIMEOUT+1)+MIN_TIMEOUT)) // if/as appropriate var electionOK bool rf.Mtx.Lock() //var _currentState = rf.CurrentState - var _currentTerm = rf.CurrentTerm - var _receivedTerm = _currentTerm // _receivedTerm is going to change - var _candidateID = rf.me - var _nServer = len(rf.peers) + // let us follow this convention when var is not declared + _currentTerm := rf.CurrentTerm + _receivedTerm := _currentTerm // _receivedTerm is going to change + _candidateID := rf.me + _nServer := len(rf.peers) if rf.CurrentState == STATE_FOLLOWER || rf.CurrentState == STATE_CANDIDATE { electionOK = true rf.CurrentTerm++ _currentTerm = rf.CurrentTerm rf.CurrentState = STATE_CANDIDATE - //_currentState = STATE_CANDIDATE - rf.VotedFor = -1 // TODO check + rf.VotedFor = -1 // TODO check what do we need to check? log.Println("[", rf.me, "] is starting an Election for Term ", _currentTerm, " at time = ", t2.Sub(t1)) } + rf.Mtx.Unlock() if electionOK { + // why we need a seperate go thread to start election go func(__currentTerm int, __receivedTerm int, __candidateID int, __nServer int) { var VoteCountForThisTerm int @@ -518,11 +532,11 @@ func Make(peers []*labrpc.ClientEnd, me int, // instant messages to the election go routine // rpc routines sends messages as soon as a new +1 vote // arrives - var VoteChannel = make(chan struct{}) + VoteChannel := make(chan struct{}) for serverID := 0; serverID < __nServer; serverID++ { //log.Println("sent ", serverID) - var req = RequestVoteArgs{} - var rep = RequestVoteReply{} + req := RequestVoteArgs{} + rep := RequestVoteReply{} req.CandidateID = __candidateID req.LastLogIndex = 0 req.LastLogTerm = 0 @@ -532,10 +546,10 @@ func Make(peers []*labrpc.ClientEnd, me int, //wg.Add(1) go func(_serverID int, _req RequestVoteArgs, _rep RequestVoteReply, __voteChan chan struct{}) { //defer wg.Done() - ok := rf.sendRequestVote(_serverID, &_req, &_rep) - //log.Println("RPC responded!") - //log.Println(rep) - if !ok { + ok := rf.sendRequestVote(_serverID, &_req, &_rep); + if !ok { + //log.Println("RPC responded!") + //log.Println(rep) // log.Println("Error in Sending RequestVote RPC!") } else { if _rep.Term > __currentTerm { @@ -546,6 +560,7 @@ func Make(peers []*labrpc.ClientEnd, me int, // // TODO check // rf.CurrentState = STATE_FOLLOWER // rf.Mtx.Unlock() + // why we haven't stopped the lection immediately here? localLock.Lock() __receivedTerm = max(__receivedTerm, _rep.Term) localLock.Unlock() @@ -553,7 +568,7 @@ func Make(peers []*labrpc.ClientEnd, me int, if _rep.VoteGranted { localLock.Lock() //log.Println("[", __candidateID, "]>", " got 1 vote in term ", __currentTerm) - VoteCountForThisTerm++ // TODO FIX, redundant for now + VoteCountForThisTerm++ // TODO FIX, redundant for now => why it is redundant? // as well as data race localLock.Unlock() __voteChan <- struct{}{} @@ -565,13 +580,16 @@ func Make(peers []*labrpc.ClientEnd, me int, } //wg.Wait() - var VoteCnt = 0 - var canStopWait = false + // this is a brilliant piece of engineering here + // best usecase of channels for synchronization among go goutines + voteCnt := 0 + canStopWait := false for !canStopWait { select { case <-VoteChannel: - VoteCnt++ - if VoteCnt >= _nServer/2 { + voteCnt++ + // check if we received majority of votes + if voteCnt >= _nServer/2 { canStopWait = true } } @@ -580,6 +598,7 @@ func Make(peers []*labrpc.ClientEnd, me int, // log.Println(__candidateID, "> Waiting done for Term ", __currentTerm, " election to over") rf.Mtx.Lock() // log.Println("[", rf.me, "]> my current term", rf.CurrentTerm) + // if some other server has term greater than us, then we convert to follower if __receivedTerm > __currentTerm { rf.CurrentState = STATE_FOLLOWER rf.CurrentTerm = __receivedTerm @@ -589,7 +608,7 @@ func Make(peers []*labrpc.ClientEnd, me int, if VoteCountForThisTerm > (len(rf.peers) / 2) { rf.CurrentState = STATE_LEADER log.Println("[", rf.me, "] : is the Leader now for Term", rf.CurrentTerm) - // Send messages to AppendEntried go routine to send + // Send messages to AppendEntries go routine to send // heartbeats to the followers // that go routine will sleep for some time and wakes up // check this state @@ -597,6 +616,7 @@ func Make(peers []*labrpc.ClientEnd, me int, rf.AppendEChannel <- struct{}{} } } else { + //a new election might have started of the server might have been moved to a new term log.Println(__candidateID, "> Something might have happened in between!") } } @@ -604,13 +624,14 @@ func Make(peers []*labrpc.ClientEnd, me int, }(_currentTerm, _receivedTerm, _candidateID, _nServer) } case <-rf.ElecTimeOutChannel: + // if we recive heartbeat message from other leader we cancel the current election and set a new timer // got message - cancel existing ElectionTimer and get new one // t2 := time.Now() // log.Println("Something Received .................... Heartbeats .....") cancel() d, cancel = context.WithTimeout(context.Background(), time.Millisecond*time.Duration(rand.Intn(MAX_TIMEOUT-MIN_TIMEOUT+1)+MIN_TIMEOUT)) // fmt.Println("Election Timer Reset after ", t2.Sub(t1)) - } + } // select closed } }() @@ -620,9 +641,10 @@ func Make(peers []*labrpc.ClientEnd, me int, for { d, cancel := context.WithTimeout(context.Background(), time.Millisecond*time.Duration(50)) select { + case <-d.Done(): // Append Entries TimeOut is over - var isLeader = false + isLeader := false rf.Mtx.Lock() if rf.CurrentState == STATE_LEADER { isLeader = true @@ -630,34 +652,34 @@ func Make(peers []*labrpc.ClientEnd, me int, rf.Mtx.Unlock() if isLeader { rf.Mtx.Lock() - //var _currentState = rf.CurrentState - var _currentTerm = rf.CurrentTerm - var _nServer = len(rf.peers) - var _prevLogIndex = 0 // TODO FIX it later - var _prevLogTerm = 0 // TODO FIX it later - var me = rf.me + _currentTerm := rf.CurrentTerm + _nServer := len(rf.peers) + _prevLogIndex := 0 // TODO FIX it later + _prevLogTerm := 0 // TODO FIX it later + me := rf.me rf.Mtx.Unlock() go sendAppendEntries(rf, _nServer, _currentTerm, _prevLogIndex, _prevLogTerm, me) } case <-rf.AppendEChannel: + // cancel exiting timer and get new one since the server just became leader cancel() d, cancel = context.WithTimeout(context.Background(), time.Millisecond*time.Duration(50)) // This go routine will hanadle sending of Empty AppendEntries Messages // to the peers //_AppendEntries(rf *Raft, _nServer int, _currentTerm int, _prevLogIndex int,_prevLogTerm int, me int) rf.Mtx.Lock() - //var _currentState = rf.CurrentState - var _currentTerm = rf.CurrentTerm - var _nServer = len(rf.peers) - var _prevLogIndex = 0 // TODO FIX it later - var _prevLogTerm = 0 // TODO FIX it later - var me = rf.me + _currentTerm := rf.CurrentTerm + _nServer := len(rf.peers) + _prevLogIndex := 0 // TODO FIX it later + _prevLogTerm := 0 // TODO FIX it later + me := rf.me rf.Mtx.Unlock() go sendAppendEntries(rf, _nServer, _currentTerm, _prevLogIndex, _prevLogTerm, me) - } - } + + } // select close + } // for close }() return rf