Skip to content
File

Blob: chord/local_chord.go

go304 lines
1package chord
2 
3import (
4 "context"
5 "fmt"
6 
7 "go.miragespace.co/specter/spec/chord"
8 "go.miragespace.co/specter/spec/protocol"
9 
10 "go.uber.org/zap"
11)
12 
13func (n *LocalNode) ID() uint64 {
14 return n.NodeConfig.Identity.GetId()
15}
16 
17func (n *LocalNode) Identity() *protocol.Node {
18 return n.NodeConfig.Identity
19}
20func (n *LocalNode) checkNodeState(leavingIsError bool) error {
21 state := n.state.Get()
22 switch state {
23 case chord.Inactive:
24 return chord.ErrNodeNotStarted
25 case chord.Leaving:
26 // get around during leaving routine, successor will ping us
27 // before replacing their predecessor pointer to our predecessor
28 if !leavingIsError {
29 return nil
30 }
31 return chord.ErrNodeGone
32 case chord.Left:
33 return chord.ErrNodeGone
34 default:
35 return nil
36 }
37}
38 
39func (n *LocalNode) Ping() error {
40 return n.checkNodeState(true)
41}
42 
43func (n *LocalNode) Notify(predecessor chord.VNode) error {
44 if err := n.checkNodeState(false); err != nil {
45 return err
46 }
47 
48 var (
49 surrogateSnapshot chord.VNode
50 predecessorSnapshot chord.VNode
51 candidatePredecessor chord.VNode
52 )
53 
54 defer func() {
55 if candidatePredecessor == nil {
56 return
57 }
58 n.surrogateMu.Lock()
59 if surrogateSnapshot == n.surrogate {
60 if candidatePredecessor.ID() == n.ID() {
61 n.surrogate = nil
62 } else {
63 n.surrogate = candidatePredecessor
64 }
65 }
66 n.surrogateMu.Unlock()
67 
68 n.predecessorMu.Lock()
69 if predecessorSnapshot == n.predecessor {
70 n.predecessor = candidatePredecessor
71 }
72 n.predecessorMu.Unlock()
73 }()
74 
75 n.surrogateMu.RLock()
76 surrogateSnapshot = n.surrogate
77 n.surrogateMu.RUnlock()
78 
79 n.predecessorMu.RLock()
80 predecessorSnapshot = n.predecessor
81 n.predecessorMu.RUnlock()
82 
83 // predecessor has not changed
84 if predecessorSnapshot != nil && predecessorSnapshot.ID() == predecessor.ID() {
85 return nil
86 }
87 
88 // we have no predecessor
89 if predecessorSnapshot == nil {
90 candidatePredecessor = predecessor
91 n.logger.Info("Discovered new predecessor via Notify",
92 zap.String("previous", "nil"),
93 zap.Object("predecessor", predecessor.Identity()),
94 )
95 return nil
96 }
97 
98 // new predecessor, check connectivity
99 if err := predecessorSnapshot.Ping(); err == nil {
100 // old predecessor is still alive, check if new predecessor is "closer"
101 if chord.Between(predecessorSnapshot.ID(), predecessor.ID(), n.ID(), false) {
102 candidatePredecessor = predecessor
103 }
104 } else {
105 // old predecessor is dead
106 candidatePredecessor = predecessor
107 }
108 
109 if candidatePredecessor != nil {
110 n.logger.Info("Discovered new predecessor via Notify",
111 zap.Object("previous", predecessorSnapshot.Identity()),
112 zap.Object("predecessor", candidatePredecessor.Identity()),
113 )
114 }
115 
116 return nil
117}
118 
119func (n *LocalNode) getSuccessor() chord.VNode {
120 n.successorsMu.RLock()
121 s := n.successors
122 n.successorsMu.RUnlock()
123 if len(s) == 0 {
124 return nil
125 }
126 return s[0]
127}
128 
129func (n *LocalNode) FindSuccessor(key uint64) (chord.VNode, error) {
130 if err := n.checkNodeState(false); err != nil {
131 return nil, err
132 }
133 pre := n.getPredecessor()
134 if pre != nil && chord.Between(pre.ID(), key, n.ID(), true) {
135 return n, nil
136 }
137 succ := n.getSuccessor()
138 if succ == nil {
139 return nil, chord.ErrNodeNoSuccessor
140 }
141 // immediate successor
142 if chord.Between(n.ID(), key, succ.ID(), true) {
143 return succ, nil
144 }
145 // find next in ring according to finger table
146 closest := n.closestPrecedingNode(key)
147 // contact possibly remote node
148 return closest.FindSuccessor(key)
149}
150 
151func (n *LocalNode) fingerRangeView(fn func(k int, f chord.VNode) bool) {
152 done := false
153 for k := chord.MaxFingerEntries; k >= 1; k-- {
154 if done {
155 return
156 }
157 n.fingers[k].computeView(func(node chord.VNode) {
158 if node != nil {
159 if !fn(k, node) {
160 done = true
161 }
162 }
163 })
164 }
165}
166 
167func (n *LocalNode) closestPrecedingNode(key uint64) chord.VNode {
168 var finger chord.VNode
169 n.fingerRangeView(func(_ int, f chord.VNode) bool {
170 if chord.Between(n.ID(), f.ID(), key, false) {
171 finger = f
172 return false
173 }
174 return true
175 })
176 if finger != nil {
177 return finger
178 }
179 // fallback to ourselves
180 return n
181}
182 
183func (n *LocalNode) getSuccessors() []chord.VNode {
184 n.successorsMu.RLock()
185 s := n.successors
186 n.successorsMu.RUnlock()
187 
188 list := make([]chord.VNode, 0, chord.ExtendedSuccessorEntries)
189 
190 for _, s := range s {
191 if s == nil {
192 continue
193 }
194 list = append(list, s)
195 }
196 
197 return list
198}
199 
200func (n *LocalNode) GetSuccessors() ([]chord.VNode, error) {
201 if err := n.checkNodeState(false); err != nil {
202 return nil, err
203 }
204 
205 return n.getSuccessors(), nil
206}
207 
208func (n *LocalNode) getPredecessor() chord.VNode {
209 n.predecessorMu.RLock()
210 p := n.predecessor
211 n.predecessorMu.RUnlock()
212 return p
213}
214 
215func (n *LocalNode) GetPredecessor() (chord.VNode, error) {
216 if err := n.checkNodeState(false); err != nil {
217 return nil, err
218 }
219 return n.getPredecessor(), nil
220}
221 
222// transferKeyUpward is called when a new predecessor has notified us, and we should transfer predecessors'
223// key range (upward). Caller of this function should hold the surrogateMu Write Lock.
224func (n *LocalNode) transferKeysUpward(ctx context.Context, prevPredecessor, newPredecessor chord.VNode) (err error) {
225 var (
226 keys [][]byte
227 values []*protocol.KVTransfer
228 low chord.VNode
229 )
230 
231 if newPredecessor.ID() == n.ID() {
232 return nil
233 }
234 
235 if prevPredecessor == nil {
236 low = n
237 } else {
238 low = prevPredecessor
239 }
240 
241 if !chord.Between(low.ID(), newPredecessor.ID(), n.ID(), false) {
242 n.logger.Debug("Skip transferring keys to predecessor because predecessor left", zap.Object("prev", low.Identity()), zap.Object("new", newPredecessor.Identity()))
243 return
244 }
245 
246 keys, err = n.kv.RangeKeys(ctx, low.ID(), newPredecessor.ID())
247 if err != nil {
248 return
249 }
250 if len(keys) == 0 {
251 return nil
252 }
253 
254 n.logger.Info("Transferring keys to new predecessor", zap.Object("predecessor", newPredecessor.Identity()), zap.Int("num_keys", len(keys)))
255 
256 values, err = n.kv.Export(ctx, keys)
257 if err != nil {
258 return
259 }
260 
261 err = newPredecessor.Import(ctx, keys, values)
262 if err != nil {
263 return
264 }
265 
266 // TODO: remove this when we implement replication
267 if err := n.kv.RemoveKeys(ctx, keys); err != nil {
268 n.logger.Error("Failed to remove keys from KV", zap.Error(err))
269 }
270 return
271}
272 
273// transferKeyDownward is called when current node is leaving the ring, and we should transfer all of our keys
274// to the successor (downward). Caller of this function should hold the surrogateMu Write Lock.
275func (n *LocalNode) transferKeysDownward(ctx context.Context, successor chord.VNode) error {
276 keys, err := n.kv.RangeKeys(ctx, 0, 0)
277 if err != nil {
278 return err
279 }
280 
281 if len(keys) == 0 {
282 return nil
283 }
284 
285 n.logger.Info("Transferring keys to successor", zap.Object("successor", successor.Identity()), zap.Int("num_keys", len(keys)))
286 
287 values, err := n.kv.Export(ctx, keys)
288 if err != nil {
289 return err
290 }
291 
292 // TODO: split into batches
293 if err := successor.Import(ctx, keys, values); err != nil {
294 return fmt.Errorf("storing KV to successor: %w", err)
295 }
296 
297 // TODO: remove this when we implement replication
298 if err := n.kv.RemoveKeys(ctx, keys); err != nil {
299 n.logger.Error("Failed to remove keys from KV", zap.Error(err))
300 }
301 
302 return nil
303}