File
Blob: chord/local_membership.go
| 1 | package chord |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "fmt" |
| 6 | |
| 7 | "go.miragespace.co/specter/spec/chord" |
| 8 | |
| 9 | "github.com/avast/retry-go/v5" |
| 10 | "go.uber.org/zap" |
| 11 | ) |
| 12 | |
| 13 | const ( |
| 14 | maxAttempts = 10 |
| 15 | ) |
| 16 | |
| 17 | func (n *LocalNode) Create() error { |
| 18 | if _, ok := n.state.Transition(chord.Inactive, chord.Joining); !ok { |
| 19 | return fmt.Errorf("node is not Inactive") |
| 20 | } |
| 21 | |
| 22 | n.logger.Info("Creating new Chord ring") |
| 23 | |
| 24 | successors := chord.MakeSuccListByID(n, []chord.VNode{}, chord.ExtendedSuccessorEntries) |
| 25 | n.successorsMu.Lock() |
| 26 | n.succListHash.Store(n.hash(successors)) |
| 27 | n.successors = successors |
| 28 | n.successorsMu.Unlock() |
| 29 | |
| 30 | for i := 1; i <= chord.MaxFingerEntries; i++ { |
| 31 | n.fingers[i].computeUpdate(func(entry *fingerEntry) { |
| 32 | entry.node = n |
| 33 | }) |
| 34 | } |
| 35 | |
| 36 | n.startTasks() |
| 37 | |
| 38 | n.state.Set(chord.Active) |
| 39 | |
| 40 | return nil |
| 41 | } |
| 42 | |
| 43 | func (n *LocalNode) Join(peer chord.VNode) error { |
| 44 | if _, ok := n.state.Transition(chord.Inactive, chord.Joining); !ok { |
| 45 | return fmt.Errorf("node is not Inactive") |
| 46 | } |
| 47 | |
| 48 | predecessor, successors, err := n.executeJoin(peer) |
| 49 | if err != nil { |
| 50 | n.state.Set(chord.Inactive) |
| 51 | return err |
| 52 | } |
| 53 | |
| 54 | n.successorsMu.Lock() |
| 55 | n.succListHash.Store(n.hash(successors)) |
| 56 | n.successors = successors |
| 57 | n.successorsMu.Unlock() |
| 58 | |
| 59 | n.predecessorMu.Lock() |
| 60 | n.predecessor = predecessor |
| 61 | n.predecessorMu.Unlock() |
| 62 | |
| 63 | n.startTasks() |
| 64 | |
| 65 | n.logger.Info("Successfully joined Chord ring", zap.Object("predecessor", predecessor.Identity()), zap.Object("successor", successors[0].Identity())) |
| 66 | |
| 67 | if err := predecessor.FinishJoin(true, false); err != nil { // advisory to let predecessor update successor list |
| 68 | n.logger.Warn("error sending advisory to predecessor", zap.Error(err)) |
| 69 | } |
| 70 | n.state.Set(chord.Active) // release local join lock |
| 71 | if err := successors[0].FinishJoin(false, true); err != nil { // release successor join lock |
| 72 | n.logger.Warn("error releasing join lock in successor", zap.Error(err)) |
| 73 | } |
| 74 | |
| 75 | return nil |
| 76 | } |
| 77 | |
| 78 | func (n *LocalNode) executeJoin(peer chord.VNode) (predecessor chord.VNode, successors []chord.VNode, err error) { |
| 79 | retrier := retry.New( |
| 80 | retry.Attempts(maxAttempts), |
| 81 | retry.Delay(n.StabilizeInterval), |
| 82 | retry.LastErrorOnly(true), |
| 83 | retry.RetryIf(chord.ErrorIsRetryable), |
| 84 | retry.OnRetry(func(attempt uint, err error) { |
| 85 | n.logger.Warn("Retrying on join error", zap.Uint("attempt", attempt), zap.Error(err)) |
| 86 | }), |
| 87 | ) |
| 88 | err = retrier.Do(func() error { |
| 89 | var joinErr error |
| 90 | n.logger.Info("Joining Chord ring", |
| 91 | zap.Object("via", peer.Identity()), |
| 92 | ) |
| 93 | predecessor, successors, joinErr = peer.RequestToJoin(n) |
| 94 | return joinErr |
| 95 | }) |
| 96 | return |
| 97 | } |
| 98 | |
| 99 | func (n *LocalNode) RequestToJoin(joiner chord.VNode) (chord.VNode, []chord.VNode, error) { |
| 100 | succ, err := n.FindSuccessor(joiner.ID()) |
| 101 | if err != nil { |
| 102 | return nil, nil, err |
| 103 | } |
| 104 | if succ.ID() == joiner.ID() { |
| 105 | return nil, nil, chord.ErrDuplicateJoinerID |
| 106 | } |
| 107 | if succ.ID() != n.ID() { |
| 108 | return succ.RequestToJoin(joiner) |
| 109 | } |
| 110 | |
| 111 | var ( |
| 112 | prevPredecessor chord.VNode |
| 113 | joined bool |
| 114 | ) |
| 115 | |
| 116 | n.logger.Info("Incoming join request", zap.Object("joiner", joiner.Identity())) |
| 117 | |
| 118 | n.surrogateMu.Lock() |
| 119 | defer n.surrogateMu.Unlock() |
| 120 | |
| 121 | // change status to transferring (if allowed) |
| 122 | if curr, ok := n.state.Transition(chord.Active, chord.Transferring); !ok { |
| 123 | n.logger.Info("Rejecting join request because current state is not Active", |
| 124 | zap.String("current", curr.String()), |
| 125 | zap.Object("joiner", joiner.Identity()), |
| 126 | ) |
| 127 | return nil, nil, chord.ErrJoinInvalidState |
| 128 | } |
| 129 | |
| 130 | // TODO: instrument how long it took to grab the lock and the duration it was held for |
| 131 | n.predecessorMu.Lock() |
| 132 | defer n.predecessorMu.Unlock() |
| 133 | |
| 134 | defer func() { |
| 135 | if joined { |
| 136 | n.predecessor = joiner |
| 137 | return |
| 138 | } |
| 139 | // joiner will request to transition. |
| 140 | // this is the reverting step upon error |
| 141 | n.state.Set(chord.Active) |
| 142 | }() |
| 143 | |
| 144 | prevPredecessor = n.predecessor |
| 145 | if prevPredecessor == nil { |
| 146 | // Without a predecessor, the key range owned by this node is unknown. |
| 147 | // Let stabilization restore it before accepting a retry of the join. |
| 148 | n.logger.Info("Rejecting join request because predecessor is unknown", |
| 149 | zap.Object("joiner", joiner.Identity()), |
| 150 | ) |
| 151 | return nil, nil, chord.ErrJoinInvalidState |
| 152 | } |
| 153 | |
| 154 | // see issue https://github.com/zllovesuki/specter/issues/23 |
| 155 | if !chord.Between(prevPredecessor.ID(), joiner.ID(), n.ID(), false) { |
| 156 | return nil, nil, chord.ErrJoinInvalidSuccessor |
| 157 | } |
| 158 | |
| 159 | // transfer key range to new node, and set surrogate pointer to new node. |
| 160 | // paper calls for forwarding but that's too hard |
| 161 | // let the caller retries |
| 162 | ctx, cancel := context.WithCancel(context.Background()) |
| 163 | defer cancel() |
| 164 | |
| 165 | if err := n.transferKeysUpward(ctx, prevPredecessor, joiner); err != nil { |
| 166 | return nil, nil, chord.ErrJoinTransferFailure |
| 167 | } |
| 168 | joined = true |
| 169 | n.surrogate = joiner |
| 170 | |
| 171 | return prevPredecessor, chord.MakeSuccListByID(n, n.getSuccessors(), chord.ExtendedSuccessorEntries), nil |
| 172 | } |
| 173 | |
| 174 | func (n *LocalNode) FinishJoin(stabilize bool, release bool) error { |
| 175 | if stabilize { |
| 176 | n.logger.Info("Join completed, joiner has requested to update pointers") |
| 177 | n.stabilize() |
| 178 | n.fixFinger() |
| 179 | } |
| 180 | if release { |
| 181 | n.logger.Info("Join completed, joiner has requested to release membership lock") |
| 182 | if curr, ok := n.state.Transition(chord.Transferring, chord.Active); !ok { |
| 183 | n.logger.Error("Unable to release membership lock", zap.String("state", curr.String())) |
| 184 | return chord.ErrJoinInvalidState |
| 185 | } |
| 186 | } |
| 187 | return nil |
| 188 | } |
| 189 | |
| 190 | func (n *LocalNode) RequestToLeave(leaver chord.VNode) error { |
| 191 | n.logger.Info("incoming leave request", zap.Object("leaver", leaver.Identity())) |
| 192 | |
| 193 | if curr, ok := n.state.Transition(chord.Active, chord.Transferring); !ok { |
| 194 | n.logger.Warn("Rejecting leave request because current state is not Active", zap.String("state", curr.String())) |
| 195 | return chord.ErrLeaveInvalidState |
| 196 | } |
| 197 | return nil |
| 198 | } |
| 199 | |
| 200 | func (n *LocalNode) FinishLeave(stabilize bool, release bool) error { |
| 201 | if stabilize { |
| 202 | n.logger.Info("Leave completed, leaver has requested to update pointers") |
| 203 | n.stabilize() |
| 204 | n.fixFinger() |
| 205 | } |
| 206 | if release { |
| 207 | n.logger.Info("Leave completed, leaver has requested to release membership lock") |
| 208 | if curr, ok := n.state.Transition(chord.Transferring, chord.Active); !ok { |
| 209 | n.logger.Error("Unable to release membership lock", zap.String("state", curr.String())) |
| 210 | return chord.ErrLeaveInvalidState |
| 211 | } |
| 212 | } |
| 213 | return nil |
| 214 | } |
| 215 | |
| 216 | func (n *LocalNode) Leave() { |
| 217 | switch n.state.Get() { |
| 218 | case chord.Inactive, chord.Leaving, chord.Left: |
| 219 | return |
| 220 | } |
| 221 | |
| 222 | left := false |
| 223 | defer func() { |
| 224 | if !left { |
| 225 | return |
| 226 | } |
| 227 | close(n.stopCh) |
| 228 | n.stopWg.Wait() |
| 229 | }() |
| 230 | |
| 231 | n.logger.Info("Requesting to leave chord ring") |
| 232 | |
| 233 | var ( |
| 234 | pre chord.VNode |
| 235 | succ chord.VNode |
| 236 | ) |
| 237 | retrier := retry.New( |
| 238 | retry.Attempts(maxAttempts), |
| 239 | retry.Delay(n.StabilizeInterval), |
| 240 | retry.LastErrorOnly(true), |
| 241 | retry.OnRetry(func(attempt uint, err error) { |
| 242 | n.logger.Warn("Retrying on leave error", zap.Uint("attempt", attempt), zap.Error(err)) |
| 243 | }), |
| 244 | ) |
| 245 | err := retrier.Do(func() error { |
| 246 | var leaveErr error |
| 247 | pre, succ, leaveErr = n.executeLeave() |
| 248 | return leaveErr |
| 249 | }) |
| 250 | if err != nil { |
| 251 | n.logger.Error("Unable to leave ring: out of attempts", zap.Error(err)) |
| 252 | return |
| 253 | } |
| 254 | |
| 255 | left = true |
| 256 | |
| 257 | n.logger.Info("Sending advisory to update pointers and releasing membership locks") |
| 258 | |
| 259 | // release membership locks |
| 260 | if pre != nil && pre.ID() != n.ID() { |
| 261 | if err := pre.FinishLeave(true, false); err != nil { // advisory to let predecessor update successor list |
| 262 | n.logger.Warn("error sending advisory to predecessor", zap.Error(err)) |
| 263 | } |
| 264 | } |
| 265 | n.state.Set(chord.Left) // release local leave lock |
| 266 | if succ != nil && succ.ID() != n.ID() { |
| 267 | if err := succ.FinishLeave(false, true); err != nil { // if applicable, release successor leave lock |
| 268 | n.logger.Warn("error releasing leave lock in successor", zap.Error(err)) |
| 269 | } |
| 270 | } |
| 271 | } |
| 272 | |
| 273 | func (n *LocalNode) executeLeave() (pre, succ chord.VNode, err error) { |
| 274 | pre = n.getPredecessor() |
| 275 | if pre == nil { |
| 276 | return nil, nil, fmt.Errorf("retrying on nil predecessor") |
| 277 | } |
| 278 | succ = n.getSuccessor() |
| 279 | if succ == nil { |
| 280 | return nil, nil, chord.ErrNodeNoSuccessor |
| 281 | } |
| 282 | if pre.ID() == n.ID() && succ.ID() == n.ID() { |
| 283 | n.logger.Debug("Skipping key transfer to successor because we are the only one left") |
| 284 | return |
| 285 | } |
| 286 | |
| 287 | // paper calls for asymmetric locking, and release locks when we are retrying |
| 288 | if n.ID() > succ.ID() { |
| 289 | if err = succ.RequestToLeave(n); err != nil { |
| 290 | return |
| 291 | } |
| 292 | if curr, ok := n.state.Transition(chord.Active, chord.Leaving); !ok { |
| 293 | n.logger.Warn("Unable to acquire local leave lock", zap.String("state", curr.String())) |
| 294 | if err := succ.FinishLeave(false, true); err != nil { // release successor lock and try again |
| 295 | n.logger.Warn("error releasing leave lock in successor", zap.Error(err)) |
| 296 | } |
| 297 | return nil, nil, chord.ErrLeaveInvalidState |
| 298 | } |
| 299 | n.logger.Info("Leave locks acquired (succ -> self)") |
| 300 | } else { |
| 301 | if curr, ok := n.state.Transition(chord.Active, chord.Leaving); !ok { |
| 302 | n.logger.Warn("Unable to acquire local leave lock", zap.String("state", curr.String())) |
| 303 | return nil, nil, chord.ErrLeaveInvalidState |
| 304 | } |
| 305 | if err := succ.RequestToLeave(n); err != nil { |
| 306 | n.state.Set(chord.Active) // release local lock and try again |
| 307 | return nil, nil, err |
| 308 | } |
| 309 | n.logger.Info("Leave locks acquired (self -> succ)") |
| 310 | } |
| 311 | |
| 312 | n.surrogateMu.Lock() |
| 313 | defer n.surrogateMu.Unlock() |
| 314 | |
| 315 | // kv requests are now blocked |
| 316 | ctx, cancel := context.WithCancel(context.Background()) |
| 317 | defer cancel() |
| 318 | |
| 319 | if err := n.transferKeysDownward(ctx, succ); err != nil { |
| 320 | n.logger.Error("Transferring KV to successor", zap.Object("successor", succ.Identity()), zap.Error(err)) |
| 321 | // release held lock and try again |
| 322 | n.state.Set(chord.Active) |
| 323 | if err := succ.FinishLeave(false, true); err != nil { |
| 324 | n.logger.Warn("error releasing leave lock in successor", zap.Error(err)) |
| 325 | } |
| 326 | return nil, nil, err |
| 327 | } |
| 328 | |
| 329 | n.surrogate = n |
| 330 | |
| 331 | return |
| 332 | } |