Skip to content
File

Blob: chord/local_tasks.go

go218 lines
1package chord
2 
3import (
4 "encoding/binary"
5 "fmt"
6 "time"
7 
8 "go.miragespace.co/specter/spec/chord"
9 "go.miragespace.co/specter/util"
10 
11 "github.com/zeebo/xxh3"
12 "go.uber.org/zap"
13)
14 
15func v2d(n []chord.VNode) []uint64 {
16 x := make([]uint64, 0)
17 for _, xx := range n {
18 if xx == nil {
19 continue
20 }
21 x = append(x, xx.ID())
22 }
23 return x
24}
25 
26// it is not safe to use xor under any ciscumstances as long as
27// we have the possbility of cyclical ring that will have ourself
28// in the successor list
29func (n *LocalNode) hash(nodes []chord.VNode) uint64 {
30 hasher := xxh3.New()
31 buf := make([]byte, 8)
32 for _, node := range nodes {
33 if node == nil {
34 continue
35 }
36 binary.BigEndian.PutUint64(buf, node.ID())
37 hasher.Write(buf)
38 }
39 return hasher.Sum64()
40}
41 
42// routine based on pseudo code from the paper "How to Make Chord Correct"
43func (n *LocalNode) stabilize() error {
44 succList := n.getSuccessors()
45 modified := false
46 
47 for len(succList) > 0 {
48 head := succList[0]
49 if head == nil {
50 return chord.ErrNodeNoSuccessor
51 }
52 newSucc, spErr := head.GetPredecessor()
53 newSuccList, nsErr := head.GetSuccessors()
54 if spErr == nil && nsErr == nil {
55 succList = chord.MakeSuccListByID(head, newSuccList, chord.ExtendedSuccessorEntries)
56 modified = true
57 
58 if newSucc != nil && chord.Between(n.ID(), newSucc.ID(), head.ID(), false) {
59 newSuccList, nsErr = newSucc.GetSuccessors()
60 if nsErr == nil {
61 succList = chord.MakeSuccListByID(newSucc, newSuccList, chord.ExtendedSuccessorEntries)
62 modified = true
63 }
64 }
65 break
66 }
67 n.logger.Debug("Skipping over successor", zap.Object("head", head.Identity()), zap.Uint64s("succ", v2d(succList)))
68 succList = succList[1:]
69 }
70 
71 n.lastStabilized.Store(time.Now())
72 
73 listHash := n.hash(succList)
74 if modified && n.succListHash.Load() != listHash {
75 n.successorsMu.Lock()
76 n.updateSuccessorsList(listHash, succList)
77 n.successorsMu.Unlock()
78 }
79 
80 if modified && len(succList) > 0 && n.checkNodeState(true) == nil { // don't re-notify our successor when we are leaving
81 succ := succList[0]
82 if err := succ.Notify(n); err != nil {
83 n.logger.Error("Error notifying successor about us", zap.Object("successor", succ.Identity()), zap.Error(err))
84 }
85 }
86 
87 return nil
88}
89 
90func (n *LocalNode) updateSuccessorsList(listHash uint64, succList []chord.VNode) {
91 n.succListHash.Store(listHash)
92 n.successors = succList
93 
94 n.logger.Info("Discovered new successors via Stabilize",
95 zap.Uint64s("successors", v2d(succList)),
96 )
97}
98 
99func (n *LocalNode) fixK(k int) (updated bool, err error) {
100 var f chord.VNode
101 next := chord.ModuloSum(n.ID(), 1<<(k-1))
102 f, err = n.FindSuccessor(next)
103 if err != nil {
104 return
105 }
106 if f == nil {
107 err = fmt.Errorf("no successor found for k = %d", k)
108 return
109 }
110 n.fingers[k].computeUpdate(func(entry *fingerEntry) {
111 if entry.node == nil || entry.node.ID() != f.ID() {
112 entry.node = f
113 updated = true
114 }
115 })
116 return
117}
118 
119func (n *LocalNode) fixFinger() error {
120 fixed := make([]int, 0)
121 for k := 1; k <= chord.MaxFingerEntries; k++ {
122 changed, err := n.fixK(k)
123 if err != nil {
124 continue
125 }
126 if changed {
127 fixed = append(fixed, k)
128 }
129 }
130 if len(fixed) > 0 {
131 n.logger.Info("FingerTable entries updated", zap.Ints("fixed", fixed))
132 }
133 return nil
134}
135 
136func (n *LocalNode) checkPredecessor() error {
137 n.predecessorMu.RLock()
138 pre := n.predecessor
139 n.predecessorMu.RUnlock()
140 if pre == nil || pre.ID() == n.ID() {
141 return nil
142 }
143 
144 err := pre.Ping()
145 if err != nil {
146 n.predecessorMu.Lock()
147 if n.predecessor == pre {
148 n.predecessor = nil
149 n.logger.Info("Discovered dead predecessor",
150 zap.Object("old", pre.Identity()),
151 zap.String("new", "nil"),
152 )
153 }
154 n.predecessorMu.Unlock()
155 }
156 return err
157}
158 
159func (n *LocalNode) periodicStabilize() {
160 defer n.stopWg.Done()
161 
162 time.Sleep(n.StabilizeInterval)
163 for {
164 select {
165 case <-n.stopCh:
166 n.logger.Debug("Stopping Stabilize task")
167 return
168 default:
169 if err := n.stabilize(); err != nil {
170 n.logger.Error("Stabilize task", zap.Error(err))
171 }
172 time.Sleep(util.RandomTimeRange(n.StabilizeInterval))
173 }
174 }
175}
176 
177func (n *LocalNode) periodicPredecessorCheck() {
178 defer n.stopWg.Done()
179 
180 for {
181 select {
182 case <-n.stopCh:
183 n.logger.Debug("Stopping predecessor checking task")
184 return
185 default:
186 n.checkPredecessor()
187 time.Sleep(util.RandomTimeRange(n.PredecessorCheckInterval))
188 }
189 }
190}
191 
192func (n *LocalNode) periodicFixFingers() {
193 defer n.stopWg.Done()
194 
195 time.Sleep(n.FixFingerInterval)
196 for {
197 select {
198 case <-n.stopCh:
199 n.logger.Debug("Stopping FixFinger task")
200 return
201 default:
202 n.fixFinger()
203 time.Sleep(util.RandomTimeRange(n.FixFingerInterval))
204 }
205 }
206}
207 
208func (n *LocalNode) startTasks() {
209 // run once
210 n.stabilize()
211 n.fixFinger()
212 n.stopWg.Add(3)
213 // then run periodically
214 go n.periodicStabilize()
215 go n.periodicPredecessorCheck()
216 go n.periodicFixFingers()
217}