Browse Source

raft: sm -> r

Blake Mizerany 11 years ago
parent
commit
79de1dc339
1 changed files with 201 additions and 198 deletions
  1. 201 198
      raft2/raft.go

+ 201 - 198
raft2/raft.go

@@ -156,12 +156,12 @@ func newStateMachine(id int64, peers []int64) *raft {
 	if id == none {
 		panic("cannot use none id")
 	}
-	sm := &raft{id: id, lead: none, raftLog: newLog(), ins: make(map[int64]*progress)}
+	r := &raft{id: id, lead: none, raftLog: newLog(), ins: make(map[int64]*progress)}
 	for _, p := range peers {
-		sm.ins[p] = &progress{}
+		r.ins[p] = &progress{}
 	}
-	sm.reset(0)
-	return sm
+	r.reset(0)
+	return r
 }
 
 func (r *raft) hasLeader() bool { return r.state != stateCandidate }
@@ -170,24 +170,24 @@ func (r *raft) propose(data []byte) {
 	r.Step(Message{From: r.id, Type: msgProp, Entries: []Entry{{Data: data}}})
 }
 
-func (sm *raft) String() string {
-	s := fmt.Sprintf(`state=%v term=%d`, sm.state, sm.Term)
-	switch sm.state {
+func (r *raft) String() string {
+	s := fmt.Sprintf(`state=%v term=%d`, r.state, r.Term)
+	switch r.state {
 	case stateFollower:
-		s += fmt.Sprintf(" vote=%v lead=%v", sm.Vote, sm.lead)
+		s += fmt.Sprintf(" vote=%v lead=%v", r.Vote, r.lead)
 	case stateCandidate:
-		s += fmt.Sprintf(` votes="%v"`, sm.votes)
+		s += fmt.Sprintf(` votes="%v"`, r.votes)
 	case stateLeader:
-		s += fmt.Sprintf(` ins="%v"`, sm.ins)
+		s += fmt.Sprintf(` ins="%v"`, r.ins)
 	}
 	return s
 }
 
-func (sm *raft) poll(id int64, v bool) (granted int) {
-	if _, ok := sm.votes[id]; !ok {
-		sm.votes[id] = v
+func (r *raft) poll(id int64, v bool) (granted int) {
+	if _, ok := r.votes[id]; !ok {
+		r.votes[id] = v
 	}
-	for _, vv := range sm.votes {
+	for _, vv := range r.votes {
 		if vv {
 			granted++
 		}
@@ -196,325 +196,328 @@ func (sm *raft) poll(id int64, v bool) (granted int) {
 }
 
 // send persists state to stable storage and then sends to its mailbox.
-func (sm *raft) send(m Message) {
-	m.From = sm.id
-	m.Term = sm.Term
-	sm.msgs = append(sm.msgs, m)
+func (r *raft) send(m Message) {
+	m.From = r.id
+	m.Term = r.Term
+	r.msgs = append(r.msgs, m)
 }
 
 // sendAppend sends RRPC, with entries to the given peer.
-func (sm *raft) sendAppend(to int64) {
-	in := sm.ins[to]
+func (r *raft) sendAppend(to int64) {
+	in := r.ins[to]
 	m := Message{}
 	m.To = to
 	m.Index = in.next - 1
-	if sm.needSnapshot(m.Index) {
+	if r.needSnapshot(m.Index) {
 		m.Type = msgSnap
-		m.Snapshot = sm.raftLog.snapshot
+		m.Snapshot = r.raftLog.snapshot
 	} else {
 		m.Type = msgApp
-		m.LogTerm = sm.raftLog.term(in.next - 1)
-		m.Entries = sm.raftLog.entries(in.next)
-		m.Commit = sm.raftLog.committed
+		m.LogTerm = r.raftLog.term(in.next - 1)
+		m.Entries = r.raftLog.entries(in.next)
+		m.Commit = r.raftLog.committed
 	}
-	sm.send(m)
+	r.send(m)
 }
 
 // sendHeartbeat sends RRPC, without entries to the given peer.
-func (sm *raft) sendHeartbeat(to int64) {
-	in := sm.ins[to]
-	index := max(in.next-1, sm.raftLog.lastIndex())
+func (r *raft) sendHeartbeat(to int64) {
+	in := r.ins[to]
+	index := max(in.next-1, r.raftLog.lastIndex())
 	m := Message{
 		To:      to,
 		Type:    msgApp,
 		Index:   index,
-		LogTerm: sm.raftLog.term(index),
-		Commit:  sm.raftLog.committed,
+		LogTerm: r.raftLog.term(index),
+		Commit:  r.raftLog.committed,
 	}
-	sm.send(m)
+	r.send(m)
 }
 
-// bcastAppend sends RRPC, with entries to all peers that are not up-to-date according to sm.mis.
-func (sm *raft) bcastAppend() {
-	for i := range sm.ins {
-		if i == sm.id {
+// bcastAppend sends RRPC, with entries to all peers that are not up-to-date according to r.mis.
+func (r *raft) bcastAppend() {
+	for i := range r.ins {
+		if i == r.id {
 			continue
 		}
-		sm.sendAppend(i)
+		r.sendAppend(i)
 	}
 }
 
 // bcastHeartbeat sends RRPC, without entries to all the peers.
-func (sm *raft) bcastHeartbeat() {
-	for i := range sm.ins {
-		if i == sm.id {
+func (r *raft) bcastHeartbeat() {
+	for i := range r.ins {
+		if i == r.id {
 			continue
 		}
-		sm.sendHeartbeat(i)
+		r.sendHeartbeat(i)
 	}
 }
 
-func (sm *raft) maybeCommit() bool {
+func (r *raft) maybeCommit() bool {
 	// TODO(bmizerany): optimize.. Currently naive
-	mis := make(int64Slice, 0, len(sm.ins))
-	for i := range sm.ins {
-		mis = append(mis, sm.ins[i].match)
+	mis := make(int64Slice, 0, len(r.ins))
+	for i := range r.ins {
+		mis = append(mis, r.ins[i].match)
 	}
 	sort.Sort(sort.Reverse(mis))
-	mci := mis[sm.q()-1]
+	mci := mis[r.q()-1]
 
-	return sm.raftLog.maybeCommit(mci, sm.Term)
+	return r.raftLog.maybeCommit(mci, r.Term)
 }
 
 // nextEnts returns the appliable entries and updates the applied index
-func (sm *raft) nextEnts() (ents []Entry) {
-	ents = sm.raftLog.nextEnts()
-	sm.raftLog.resetNextEnts()
+func (r *raft) nextEnts() (ents []Entry) {
+	ents = r.raftLog.nextEnts()
+	r.raftLog.resetNextEnts()
 	return ents
 }
 
-func (sm *raft) reset(term int64) {
-	sm.setTerm(term)
-	sm.lead.Set(none)
-	sm.setVote(none)
-	sm.votes = make(map[int64]bool)
-	for i := range sm.ins {
-		sm.ins[i] = &progress{next: sm.raftLog.lastIndex() + 1}
-		if i == sm.id {
-			sm.ins[i].match = sm.raftLog.lastIndex()
+func (r *raft) reset(term int64) {
+	r.setTerm(term)
+	r.lead.Set(none)
+	r.setVote(none)
+	r.votes = make(map[int64]bool)
+	for i := range r.ins {
+		r.ins[i] = &progress{next: r.raftLog.lastIndex() + 1}
+		if i == r.id {
+			r.ins[i].match = r.raftLog.lastIndex()
 		}
 	}
 }
 
-func (sm *raft) q() int {
-	return len(sm.ins)/2 + 1
+func (r *raft) q() int {
+	return len(r.ins)/2 + 1
 }
 
-func (sm *raft) appendEntry(e Entry) {
-	e.Term = sm.Term
-	e.Index = sm.raftLog.lastIndex() + 1
-	sm.LastIndex = sm.raftLog.append(sm.raftLog.lastIndex(), e)
-	sm.ins[sm.id].update(sm.raftLog.lastIndex())
-	sm.maybeCommit()
+func (r *raft) appendEntry(e Entry) {
+	e.Term = r.Term
+	e.Index = r.raftLog.lastIndex() + 1
+	r.LastIndex = r.raftLog.append(r.raftLog.lastIndex(), e)
+	r.ins[r.id].update(r.raftLog.lastIndex())
+	r.maybeCommit()
 }
 
-func (sm *raft) becomeFollower(term int64, lead int64) {
-	sm.reset(term)
-	sm.lead.Set(lead)
-	sm.state = stateFollower
-	sm.pendingConf = false
+func (r *raft) becomeFollower(term int64, lead int64) {
+	r.reset(term)
+	r.lead.Set(lead)
+	r.state = stateFollower
+	r.pendingConf = false
 }
 
-func (sm *raft) becomeCandidate() {
+func (r *raft) becomeCandidate() {
 	// TODO(xiangli) remove the panic when the raft implementation is stable
-	if sm.state == stateLeader {
+	if r.state == stateLeader {
 		panic("invalid transition [leader -> candidate]")
 	}
-	sm.reset(sm.Term + 1)
-	sm.setVote(sm.id)
-	sm.state = stateCandidate
+	r.reset(r.Term + 1)
+	r.setVote(r.id)
+	r.state = stateCandidate
 }
 
-func (sm *raft) becomeLeader() {
+func (r *raft) becomeLeader() {
 	// TODO(xiangli) remove the panic when the raft implementation is stable
-	if sm.state == stateFollower {
+	if r.state == stateFollower {
 		panic("invalid transition [follower -> leader]")
 	}
-	sm.reset(sm.Term)
-	sm.lead.Set(sm.id)
-	sm.state = stateLeader
+	r.reset(r.Term)
+	r.lead.Set(r.id)
+	r.state = stateLeader
 
-	for _, e := range sm.raftLog.entries(sm.raftLog.committed + 1) {
+	for _, e := range r.raftLog.entries(r.raftLog.committed + 1) {
 		if e.isConfig() {
-			sm.pendingConf = true
+			r.pendingConf = true
 		}
 	}
 
-	sm.appendEntry(Entry{Type: Normal, Data: nil})
+	r.appendEntry(Entry{Type: Normal, Data: nil})
 }
 
-func (sm *raft) ReadMessages() []Message {
-	msgs := sm.msgs
-	sm.msgs = make([]Message, 0)
+func (r *raft) ReadMessages() []Message {
+	msgs := r.msgs
+	r.msgs = make([]Message, 0)
 
 	return msgs
 }
 
-func (sm *raft) Step(m Message) error {
+func (r *raft) Step(m Message) error {
+	// TODO(bmizerany): this likely allocs - prevent that.
+	defer func() { r.Commit = r.raftLog.committed }()
+
 	if m.Type == msgHup {
-		sm.becomeCandidate()
-		if sm.q() == sm.poll(sm.id, true) {
-			sm.becomeLeader()
+		r.becomeCandidate()
+		if r.q() == r.poll(r.id, true) {
+			r.becomeLeader()
 		}
-		for i := range sm.ins {
-			if i == sm.id {
+		for i := range r.ins {
+			if i == r.id {
 				continue
 			}
-			lasti := sm.raftLog.lastIndex()
-			sm.send(Message{To: i, Type: msgVote, Index: lasti, LogTerm: sm.raftLog.term(lasti)})
+			lasti := r.raftLog.lastIndex()
+			r.send(Message{To: i, Type: msgVote, Index: lasti, LogTerm: r.raftLog.term(lasti)})
 		}
 	}
 
 	switch {
 	case m.Term == 0:
 		// local message
-	case m.Term > sm.Term:
+	case m.Term > r.Term:
 		lead := m.From
 		if m.Type == msgVote {
 			lead = none
 		}
-		sm.becomeFollower(m.Term, lead)
-	case m.Term < sm.Term:
+		r.becomeFollower(m.Term, lead)
+	case m.Term < r.Term:
 		// ignore
 	}
 
-	stepmap[sm.state](sm, m)
+	stepmap[r.state](r, m)
 	return nil
 }
 
-func (sm *raft) handleAppendEntries(m Message) {
-	if sm.raftLog.maybeAppend(m.Index, m.LogTerm, m.Commit, m.Entries...) {
-		sm.LastIndex = sm.raftLog.lastIndex()
-		sm.send(Message{To: m.From, Type: msgAppResp, Index: sm.raftLog.lastIndex()})
+func (r *raft) handleAppendEntries(m Message) {
+	if r.raftLog.maybeAppend(m.Index, m.LogTerm, m.Commit, m.Entries...) {
+		r.LastIndex = r.raftLog.lastIndex()
+		r.send(Message{To: m.From, Type: msgAppResp, Index: r.raftLog.lastIndex()})
 	} else {
-		sm.send(Message{To: m.From, Type: msgAppResp, Index: -1})
+		r.send(Message{To: m.From, Type: msgAppResp, Index: -1})
 	}
 }
 
-func (sm *raft) handleSnapshot(m Message) {
-	if sm.restore(m.Snapshot) {
-		sm.send(Message{To: m.From, Type: msgAppResp, Index: sm.raftLog.lastIndex()})
+func (r *raft) handleSnapshot(m Message) {
+	if r.restore(m.Snapshot) {
+		r.send(Message{To: m.From, Type: msgAppResp, Index: r.raftLog.lastIndex()})
 	} else {
-		sm.send(Message{To: m.From, Type: msgAppResp, Index: sm.raftLog.committed})
+		r.send(Message{To: m.From, Type: msgAppResp, Index: r.raftLog.committed})
 	}
 }
 
-func (sm *raft) addNode(id int64) {
-	sm.addIns(id, 0, sm.raftLog.lastIndex()+1)
-	sm.pendingConf = false
-	if id == sm.id {
-		sm.promotable = true
+func (r *raft) addNode(id int64) {
+	r.addIns(id, 0, r.raftLog.lastIndex()+1)
+	r.pendingConf = false
+	if id == r.id {
+		r.promotable = true
 	}
 }
 
-func (sm *raft) removeNode(id int64) {
-	sm.deleteIns(id)
-	sm.pendingConf = false
+func (r *raft) removeNode(id int64) {
+	r.deleteIns(id)
+	r.pendingConf = false
 }
 
-type stepFunc func(sm *raft, m Message) bool
+type stepFunc func(r *raft, m Message) bool
 
-func stepLeader(sm *raft, m Message) bool {
+func stepLeader(r *raft, m Message) bool {
 	switch m.Type {
 	case msgBeat:
-		sm.bcastHeartbeat()
+		r.bcastHeartbeat()
 	case msgProp:
 		if len(m.Entries) != 1 {
 			panic("unexpected length(entries) of a msgProp")
 		}
 		e := m.Entries[0]
 		if e.isConfig() {
-			if sm.pendingConf {
+			if r.pendingConf {
 				return false
 			}
-			sm.pendingConf = true
+			r.pendingConf = true
 		}
-		sm.appendEntry(e)
-		sm.bcastAppend()
+		r.appendEntry(e)
+		r.bcastAppend()
 	case msgAppResp:
 		if m.Index < 0 {
-			sm.ins[m.From].decr()
-			sm.sendAppend(m.From)
+			r.ins[m.From].decr()
+			r.sendAppend(m.From)
 		} else {
-			sm.ins[m.From].update(m.Index)
-			if sm.maybeCommit() {
-				sm.bcastAppend()
+			r.ins[m.From].update(m.Index)
+			if r.maybeCommit() {
+				r.bcastAppend()
 			}
 		}
 	case msgVote:
-		sm.send(Message{To: m.From, Type: msgVoteResp, Index: -1})
+		r.send(Message{To: m.From, Type: msgVoteResp, Index: -1})
 	}
 	return true
 }
 
-func stepCandidate(sm *raft, m Message) bool {
+func stepCandidate(r *raft, m Message) bool {
 	switch m.Type {
 	case msgProp:
 		return false
 	case msgApp:
-		sm.becomeFollower(sm.Term, m.From)
-		sm.handleAppendEntries(m)
+		r.becomeFollower(r.Term, m.From)
+		r.handleAppendEntries(m)
 	case msgSnap:
-		sm.becomeFollower(m.Term, m.From)
-		sm.handleSnapshot(m)
+		r.becomeFollower(m.Term, m.From)
+		r.handleSnapshot(m)
 	case msgVote:
-		sm.send(Message{To: m.From, Type: msgVoteResp, Index: -1})
+		r.send(Message{To: m.From, Type: msgVoteResp, Index: -1})
 	case msgVoteResp:
-		gr := sm.poll(m.From, m.Index >= 0)
-		switch sm.q() {
+		gr := r.poll(m.From, m.Index >= 0)
+		switch r.q() {
 		case gr:
-			sm.becomeLeader()
-			sm.bcastAppend()
-		case len(sm.votes) - gr:
-			sm.becomeFollower(sm.Term, none)
+			r.becomeLeader()
+			r.bcastAppend()
+		case len(r.votes) - gr:
+			r.becomeFollower(r.Term, none)
 		}
 	}
 	return true
 }
 
-func stepFollower(sm *raft, m Message) bool {
+func stepFollower(r *raft, m Message) bool {
 	switch m.Type {
 	case msgProp:
-		if sm.lead.Get() == none {
+		if r.lead.Get() == none {
 			return false
 		}
-		m.To = sm.lead.Get()
-		sm.send(m)
+		m.To = r.lead.Get()
+		r.send(m)
 	case msgApp:
-		sm.lead.Set(m.From)
-		sm.handleAppendEntries(m)
+		r.lead.Set(m.From)
+		r.handleAppendEntries(m)
 	case msgSnap:
-		sm.handleSnapshot(m)
+		r.handleSnapshot(m)
 	case msgVote:
-		if (sm.Vote == none || sm.Vote == m.From) && sm.raftLog.isUpToDate(m.Index, m.LogTerm) {
-			sm.setVote(m.From)
-			sm.send(Message{To: m.From, Type: msgVoteResp, Index: sm.raftLog.lastIndex()})
+		if (r.Vote == none || r.Vote == m.From) && r.raftLog.isUpToDate(m.Index, m.LogTerm) {
+			r.setVote(m.From)
+			r.send(Message{To: m.From, Type: msgVoteResp, Index: r.raftLog.lastIndex()})
 		} else {
-			sm.send(Message{To: m.From, Type: msgVoteResp, Index: -1})
+			r.send(Message{To: m.From, Type: msgVoteResp, Index: -1})
 		}
 	}
 	return true
 }
 
-func (sm *raft) compact(d []byte) {
-	sm.raftLog.snap(d, sm.raftLog.applied, sm.raftLog.term(sm.raftLog.applied), sm.nodes())
-	sm.raftLog.compact(sm.raftLog.applied)
+func (r *raft) compact(d []byte) {
+	r.raftLog.snap(d, r.raftLog.applied, r.raftLog.term(r.raftLog.applied), r.nodes())
+	r.raftLog.compact(r.raftLog.applied)
 }
 
 // restore recovers the statemachine from a snapshot. It restores the log and the
 // configuration of statemachine.
-func (sm *raft) restore(s Snapshot) bool {
-	if s.Index <= sm.raftLog.committed {
+func (r *raft) restore(s Snapshot) bool {
+	if s.Index <= r.raftLog.committed {
 		return false
 	}
 
-	sm.raftLog.restore(s)
-	sm.LastIndex = sm.raftLog.lastIndex()
-	sm.ins = make(map[int64]*progress)
+	r.raftLog.restore(s)
+	r.LastIndex = r.raftLog.lastIndex()
+	r.ins = make(map[int64]*progress)
 	for _, n := range s.Nodes {
-		if n == sm.id {
-			sm.addIns(n, sm.raftLog.lastIndex(), sm.raftLog.lastIndex()+1)
+		if n == r.id {
+			r.addIns(n, r.raftLog.lastIndex(), r.raftLog.lastIndex()+1)
 		} else {
-			sm.addIns(n, 0, sm.raftLog.lastIndex()+1)
+			r.addIns(n, 0, r.raftLog.lastIndex()+1)
 		}
 	}
-	sm.pendingConf = false
+	r.pendingConf = false
 	return true
 }
 
-func (sm *raft) needSnapshot(i int64) bool {
-	if i < sm.raftLog.offset {
-		if sm.raftLog.snapshot.IsEmpty() {
+func (r *raft) needSnapshot(i int64) bool {
+	if i < r.raftLog.offset {
+		if r.raftLog.snapshot.IsEmpty() {
 			panic("need non-empty snapshot")
 		}
 		return true
@@ -522,59 +525,59 @@ func (sm *raft) needSnapshot(i int64) bool {
 	return false
 }
 
-func (sm *raft) nodes() []int64 {
-	nodes := make([]int64, 0, len(sm.ins))
-	for k := range sm.ins {
+func (r *raft) nodes() []int64 {
+	nodes := make([]int64, 0, len(r.ins))
+	for k := range r.ins {
 		nodes = append(nodes, k)
 	}
 	return nodes
 }
 
-func (sm *raft) setTerm(term int64) {
-	sm.Term = term
-	sm.saveState()
+func (r *raft) setTerm(term int64) {
+	r.Term = term
+	r.saveState()
 }
 
-func (sm *raft) setVote(vote int64) {
-	sm.Vote = vote
-	sm.saveState()
+func (r *raft) setVote(vote int64) {
+	r.Vote = vote
+	r.saveState()
 }
 
-func (sm *raft) addIns(id, match, next int64) {
-	sm.ins[id] = &progress{next: next, match: match}
-	sm.saveState()
+func (r *raft) addIns(id, match, next int64) {
+	r.ins[id] = &progress{next: next, match: match}
+	r.saveState()
 }
 
-func (sm *raft) deleteIns(id int64) {
-	delete(sm.ins, id)
-	sm.saveState()
+func (r *raft) deleteIns(id int64) {
+	delete(r.ins, id)
+	r.saveState()
 }
 
-// saveState saves the state to sm.State
+// saveState saves the state to r.State
 // When there is a term change, vote change or configuration change, raft
 // must call saveState.
-func (sm *raft) saveState() {
-	sm.setState(sm.Vote, sm.Term, sm.raftLog.committed)
+func (r *raft) saveState() {
+	r.setState(r.Vote, r.Term, r.raftLog.committed)
 }
 
-func (sm *raft) setState(vote, term, commit int64) {
-	sm.Vote = vote
-	sm.Term = term
-	sm.Commit = commit
+func (r *raft) setState(vote, term, commit int64) {
+	r.Vote = vote
+	r.Term = term
+	r.Commit = commit
 }
 
-func (sm *raft) loadEnts(ents []Entry) {
-	if !sm.raftLog.isEmpty() {
+func (r *raft) loadEnts(ents []Entry) {
+	if !r.raftLog.isEmpty() {
 		panic("cannot load entries when log is not empty")
 	}
-	sm.raftLog.append(0, ents...)
-	sm.raftLog.unstable = sm.raftLog.lastIndex() + 1
+	r.raftLog.append(0, ents...)
+	r.raftLog.unstable = r.raftLog.lastIndex() + 1
 }
 
-func (sm *raft) loadState(state State) {
-	sm.raftLog.committed = state.Commit
-	sm.setTerm(state.Term)
-	sm.setVote(state.Vote)
+func (r *raft) loadState(state State) {
+	r.raftLog.committed = state.Commit
+	r.setTerm(state.Term)
+	r.setVote(state.Vote)
 }
 
 func (s *State) IsEmpty() bool {