File
Blob: chord/local.go
| 1 | package chord |
| 2 | |
| 3 | import ( |
| 4 | "net/http" |
| 5 | "sync" |
| 6 | "time" |
| 7 | |
| 8 | "go.miragespace.co/specter/spec/chord" |
| 9 | "go.miragespace.co/specter/util/acceptor" |
| 10 | "go.miragespace.co/specter/util/ratecounter" |
| 11 | |
| 12 | "github.com/TheZeroSlave/zapsentry" |
| 13 | "go.uber.org/atomic" |
| 14 | "go.uber.org/zap" |
| 15 | ) |
| 16 | |
| 17 | type LocalNode struct { |
| 18 | logger *zap.Logger |
| 19 | predecessorMu sync.RWMutex // simple mutex surrounding operations on predecessor |
| 20 | predecessor chord.VNode // nord's immediate predecessor |
| 21 | surrogateMu sync.RWMutex // advanced mutex surrounding KV requests during Join/Leave |
| 22 | surrogate chord.VNode // node's previous predecessor, used to guard against outdated KV requests |
| 23 | successorsMu sync.RWMutex // advanced mutex surrounding successors during Join/Leave |
| 24 | successors []chord.VNode // node's extended list of successors |
| 25 | succListHash *atomic.Uint64 // simple hash on the successors to determine if they have changed |
| 26 | kv chord.KVProvider // KV backing implementation |
| 27 | chordRate *ratecounter.Rate // track how chatty incoming chord rpc is |
| 28 | kvRate *ratecounter.Rate // track how chatty incoming kv rpc is |
| 29 | kvStaleCount *atomic.Uint64 // track the number of kv stale ownership |
| 30 | rpcErrorCount *atomic.Uint64 // track the number of non-retryable rpc request errors |
| 31 | lastStabilized *atomic.Time // last stablized time |
| 32 | stopWg sync.WaitGroup // used to wait for task goroutines to be stopped |
| 33 | stopCh chan struct{} // used to signal task goroutines to stop |
| 34 | state *nodeState // a replacement of LockQueue from the paper |
| 35 | fingers [chord.MaxFingerEntries + 1]fingerEntry // finger table to provide log(N) optimization |
| 36 | rpcAcceptor *acceptor.HTTP2Acceptor // listener for handling incoming rpc request |
| 37 | rpcHandler http.Handler |
| 38 | rpcHandlerOnce sync.Once |
| 39 | NodeConfig |
| 40 | } |
| 41 | |
| 42 | var _ chord.VNode = (*LocalNode)(nil) |
| 43 | |
| 44 | type fingerEntry struct { |
| 45 | node chord.VNode |
| 46 | sync.RWMutex |
| 47 | } |
| 48 | |
| 49 | func (f *fingerEntry) computeUpdate(fn func(entry *fingerEntry)) { |
| 50 | f.Lock() |
| 51 | defer f.Unlock() |
| 52 | |
| 53 | fn(f) |
| 54 | } |
| 55 | |
| 56 | func (f *fingerEntry) computeView(fn func(node chord.VNode)) { |
| 57 | f.RLock() |
| 58 | defer f.RUnlock() |
| 59 | |
| 60 | fn(f.node) |
| 61 | } |
| 62 | |
| 63 | func NewLocalNode(conf NodeConfig) *LocalNode { |
| 64 | if err := conf.Validate(); err != nil { |
| 65 | panic(err) |
| 66 | } |
| 67 | n := &LocalNode{ |
| 68 | NodeConfig: conf, |
| 69 | logger: conf.BaseLogger.With(zapsentry.NewScope()).With(zap.String("component", "localNode"), zap.Uint64("node", conf.Identity.GetId())), |
| 70 | state: newNodeState(chord.Inactive), |
| 71 | succListHash: atomic.NewUint64(conf.Identity.GetId()), |
| 72 | kv: conf.KVProvider, |
| 73 | lastStabilized: atomic.NewTime(time.Time{}), |
| 74 | stopCh: make(chan struct{}), |
| 75 | rpcAcceptor: acceptor.NewH2Acceptor(nil), |
| 76 | chordRate: ratecounter.New(time.Second, time.Second*10), |
| 77 | kvRate: ratecounter.New(time.Second, time.Second*10), |
| 78 | kvStaleCount: atomic.NewUint64(0), |
| 79 | rpcErrorCount: atomic.NewUint64(0), |
| 80 | } |
| 81 | |
| 82 | return n |
| 83 | } |