| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345 |
- package metadata
- /*
- ACMS replicated log.
- The leader assigns sequence numbers and fans entries out; followers apply
- them and remember the highest contiguous sequence they applied. Changes
- made on a follower go to the leader (or wait in the pending queue until a
- leader is reachable). Applying is idempotent (last-writer-wins per
- record) so duplicates and reordering are harmless; the sequence numbers
- only tell a follower when it has missed something and must catch up.
- */
- import (
- "context"
- "encoding/json"
- "net/http"
- "time"
- uuid "github.com/satori/go.uuid"
- "imuslab.com/arozos/mod/cluster/acn"
- "imuslab.com/arozos/mod/cluster/membership"
- )
- const (
- pathAppend = acn.BasePath + "/meta/append"
- pathSubmit = acn.BasePath + "/meta/submit"
- pathLog = acn.BasePath + "/meta/log"
- pathSnapshot = acn.BasePath + "/meta/snapshot"
- pathLease = acn.BasePath + "/meta/lease"
- pathStat = acn.BasePath + "/meta/stat"
- logPageSize = 500
- )
- func newEntry(kind string, rec interface{}, version int64, origin string) (Entry, error) {
- js, err := json.Marshal(rec)
- if err != nil {
- return Entry{}, err
- }
- return Entry{Kind: kind, Payload: js, Version: version, Origin: origin}, nil
- }
- // apply folds one entry into the local store. Returns true when it changed
- // anything.
- func (mgr *Manager) apply(e Entry) bool {
- changed := false
- switch e.Kind {
- case KindFile:
- var rec FileRecord
- if json.Unmarshal(e.Payload, &rec) == nil {
- changed = mgr.st.putFile(&rec)
- }
- case KindVolume:
- var v Volume
- if json.Unmarshal(e.Payload, &v) == nil {
- changed = mgr.st.putVolume(&v)
- }
- case KindPolicy:
- var p Policy
- if json.Unmarshal(e.Payload, &p) == nil {
- changed = mgr.st.putPolicy(&p)
- }
- case KindJob:
- var j Job
- if json.Unmarshal(e.Payload, &j) == nil {
- changed = mgr.st.putJob(&j)
- }
- case KindSetting:
- var st Setting
- if json.Unmarshal(e.Payload, &st) == nil {
- changed = mgr.st.putSetting(&st)
- }
- }
- if changed && mgr.OnChange != nil {
- mgr.OnChange(e.Kind, e.Payload)
- }
- return changed
- }
- // dispatch replicates an entry that was already applied locally.
- func (mgr *Manager) dispatch(e Entry) {
- if mgr.IsLeader() {
- mgr.leaderAccept(e)
- return
- }
- id := uuid.NewV4().String()
- mgr.st.addPending(id, e)
- go mgr.flushPending()
- }
- // leaderAccept assigns a sequence number and pushes the entry to every peer.
- func (mgr *Manager) leaderAccept(e Entry) Entry {
- lease := mgr.st.getLease()
- assigned := mgr.st.appendLog(e, lease.Term)
- mgr.st.setApplied(assigned.Seq, lease.Term)
- mgr.mu.Lock()
- mgr.appends++
- compact := mgr.appends%1000 == 0
- mgr.mu.Unlock()
- if compact {
- mgr.st.compactLog(mgr.opt.LogKeep)
- }
- go mgr.fanout([]Entry{assigned})
- return assigned
- }
- func (mgr *Manager) onlinePeers() []string {
- out := []string{}
- for _, n := range mgr.m.NodeViews() {
- if n.Local {
- continue
- }
- if n.State == membership.StateOnline || n.State == membership.StateDegraded {
- out = append(out, n.ID)
- }
- }
- return out
- }
- func (mgr *Manager) fanout(entries []Entry) {
- for _, peer := range mgr.onlinePeers() {
- go func(id string) {
- ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
- defer cancel()
- mgr.m.Transport().DoJSON(ctx, id, http.MethodPost, pathAppend, AppendRequest{Entries: entries}, nil)
- }(peer)
- }
- }
- // followerReceive handles entries pushed by the leader.
- func (mgr *Manager) followerReceive(entries []Entry) {
- applied, term := mgr.st.getApplied()
- for _, e := range entries {
- mgr.apply(e)
- if e.Term != term {
- //New leadership term: the numbering restarted, take a snapshot later
- term = e.Term
- applied = 0
- mgr.st.setApplied(0, term)
- mgr.requestCatchup()
- }
- if e.Seq == applied+1 {
- applied = e.Seq
- mgr.st.setApplied(applied, term)
- } else if e.Seq > applied+1 {
- mgr.requestCatchup()
- }
- mgr.st.setLastSeq(e.Seq)
- }
- }
- func (mgr *Manager) requestCatchup() {
- select {
- case mgr.catchupCh <- struct{}{}:
- default:
- }
- }
- func (mgr *Manager) catchupLoop() {
- defer mgr.wg.Done()
- for {
- select {
- case <-mgr.stop:
- return
- case <-mgr.catchupCh:
- mgr.catchup()
- }
- }
- }
- // catchup pulls missing log entries (or a snapshot) from the leader.
- func (mgr *Manager) catchup() {
- if mgr.IsLeader() {
- return
- }
- leader := mgr.Leader()
- if leader == "" {
- return
- }
- for i := 0; i < 100; i++ {
- applied, term := mgr.st.getApplied()
- if applied == 0 {
- //Fresh node or fresh term: the log may have been compacted below
- //what we need, a snapshot is the only safe starting point
- mgr.snapshotFromLeader(leader)
- return
- }
- ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
- resp, err := mgr.m.Transport().Do(ctx, leader, http.MethodGet, pathLog+"?after="+uitoa(applied), nil)
- cancel()
- if err != nil {
- return
- }
- if resp.Status == http.StatusGone {
- mgr.snapshotFromLeader(leader)
- return
- }
- if err := resp.Error(); err != nil {
- return
- }
- var page LogResponse
- if json.Unmarshal(resp.Body, &page) != nil {
- return
- }
- if page.Term != term {
- mgr.snapshotFromLeader(leader)
- return
- }
- for _, e := range page.Entries {
- mgr.apply(e)
- applied = e.Seq
- }
- mgr.st.setApplied(applied, page.Term)
- mgr.st.setLastSeq(page.LastSeq)
- if len(page.Entries) < logPageSize || applied >= page.LastSeq {
- return
- }
- }
- }
- func (mgr *Manager) snapshotFromLeader(leader string) {
- ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
- defer cancel()
- var snap Snapshot
- if err := mgr.m.Transport().DoJSON(ctx, leader, http.MethodGet, pathSnapshot, nil, &snap); err != nil {
- return
- }
- mgr.applySnapshot(snap)
- }
- func (mgr *Manager) applySnapshot(snap Snapshot) {
- for i := range snap.Files {
- if mgr.st.putFile(&snap.Files[i]) && mgr.OnChange != nil {
- js, _ := json.Marshal(snap.Files[i])
- mgr.OnChange(KindFile, js)
- }
- }
- for i := range snap.Volumes {
- if mgr.st.putVolume(&snap.Volumes[i]) && mgr.OnChange != nil {
- js, _ := json.Marshal(snap.Volumes[i])
- mgr.OnChange(KindVolume, js)
- }
- }
- for i := range snap.Policies {
- if mgr.st.putPolicy(&snap.Policies[i]) && mgr.OnChange != nil {
- js, _ := json.Marshal(snap.Policies[i])
- mgr.OnChange(KindPolicy, js)
- }
- }
- for i := range snap.Jobs {
- if mgr.st.putJob(&snap.Jobs[i]) && mgr.OnChange != nil {
- js, _ := json.Marshal(snap.Jobs[i])
- mgr.OnChange(KindJob, js)
- }
- }
- for i := range snap.Settings {
- if mgr.st.putSetting(&snap.Settings[i]) && mgr.OnChange != nil {
- js, _ := json.Marshal(snap.Settings[i])
- mgr.OnChange(KindSetting, js)
- }
- }
- mgr.st.setApplied(snap.LastSeq, snap.Term)
- mgr.st.setLastSeq(snap.LastSeq)
- }
- // pendingLoop retries follower changes that no leader has accepted yet.
- func (mgr *Manager) pendingLoop() {
- defer mgr.wg.Done()
- ticker := time.NewTicker(mgr.opt.PendingRetry)
- defer ticker.Stop()
- for {
- select {
- case <-mgr.stop:
- return
- case <-ticker.C:
- mgr.flushPending()
- }
- }
- }
- // flushPending hands queued entries to the leader (or accepts them locally
- // when this node became the leader meanwhile).
- // flushPending drains the pending queue. It is called after every local
- // write, so only one call does the work at a time: a second call just asks
- // the running one for another round. Many concurrent drains each walking the
- // whole queue made the work grow with the square of the backlog.
- func (mgr *Manager) flushPending() {
- mgr.flushAgain.Store(true)
- for {
- if !mgr.flushing.CompareAndSwap(false, true) {
- return //the active drain will see flushAgain
- }
- for mgr.flushAgain.Swap(false) {
- mgr.drainPending()
- }
- mgr.flushing.Store(false)
- //A request that arrived as we were finishing gets its round now
- if !mgr.flushAgain.Load() {
- return
- }
- }
- }
- // drainPending sends every queued change to the leader, or accepts it when
- // this node is the leader.
- func (mgr *Manager) drainPending() {
- pending := mgr.st.pendingEntries()
- if len(pending) == 0 {
- return
- }
- if mgr.IsLeader() {
- for id, e := range pending {
- mgr.leaderAccept(e)
- mgr.st.removePending(id)
- }
- return
- }
- leader := mgr.Leader()
- if leader == "" {
- return
- }
- for id, e := range pending {
- ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
- err := mgr.m.Transport().DoJSON(ctx, leader, http.MethodPost, pathSubmit, e, nil)
- cancel()
- if err != nil {
- return
- }
- mgr.st.removePending(id)
- }
- }
- func uitoa(v uint64) string {
- if v == 0 {
- return "0"
- }
- buf := [20]byte{}
- i := len(buf)
- for v > 0 {
- i--
- buf[i] = byte('0' + v%10)
- v /= 10
- }
- return string(buf[i:])
- }
|