log.go 8.4 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345
  1. package metadata
  2. /*
  3. ACMS replicated log.
  4. The leader assigns sequence numbers and fans entries out; followers apply
  5. them and remember the highest contiguous sequence they applied. Changes
  6. made on a follower go to the leader (or wait in the pending queue until a
  7. leader is reachable). Applying is idempotent (last-writer-wins per
  8. record) so duplicates and reordering are harmless; the sequence numbers
  9. only tell a follower when it has missed something and must catch up.
  10. */
  11. import (
  12. "context"
  13. "encoding/json"
  14. "net/http"
  15. "time"
  16. uuid "github.com/satori/go.uuid"
  17. "imuslab.com/arozos/mod/cluster/acn"
  18. "imuslab.com/arozos/mod/cluster/membership"
  19. )
  20. const (
  21. pathAppend = acn.BasePath + "/meta/append"
  22. pathSubmit = acn.BasePath + "/meta/submit"
  23. pathLog = acn.BasePath + "/meta/log"
  24. pathSnapshot = acn.BasePath + "/meta/snapshot"
  25. pathLease = acn.BasePath + "/meta/lease"
  26. pathStat = acn.BasePath + "/meta/stat"
  27. logPageSize = 500
  28. )
  29. func newEntry(kind string, rec interface{}, version int64, origin string) (Entry, error) {
  30. js, err := json.Marshal(rec)
  31. if err != nil {
  32. return Entry{}, err
  33. }
  34. return Entry{Kind: kind, Payload: js, Version: version, Origin: origin}, nil
  35. }
  36. // apply folds one entry into the local store. Returns true when it changed
  37. // anything.
  38. func (mgr *Manager) apply(e Entry) bool {
  39. changed := false
  40. switch e.Kind {
  41. case KindFile:
  42. var rec FileRecord
  43. if json.Unmarshal(e.Payload, &rec) == nil {
  44. changed = mgr.st.putFile(&rec)
  45. }
  46. case KindVolume:
  47. var v Volume
  48. if json.Unmarshal(e.Payload, &v) == nil {
  49. changed = mgr.st.putVolume(&v)
  50. }
  51. case KindPolicy:
  52. var p Policy
  53. if json.Unmarshal(e.Payload, &p) == nil {
  54. changed = mgr.st.putPolicy(&p)
  55. }
  56. case KindJob:
  57. var j Job
  58. if json.Unmarshal(e.Payload, &j) == nil {
  59. changed = mgr.st.putJob(&j)
  60. }
  61. case KindSetting:
  62. var st Setting
  63. if json.Unmarshal(e.Payload, &st) == nil {
  64. changed = mgr.st.putSetting(&st)
  65. }
  66. }
  67. if changed && mgr.OnChange != nil {
  68. mgr.OnChange(e.Kind, e.Payload)
  69. }
  70. return changed
  71. }
  72. // dispatch replicates an entry that was already applied locally.
  73. func (mgr *Manager) dispatch(e Entry) {
  74. if mgr.IsLeader() {
  75. mgr.leaderAccept(e)
  76. return
  77. }
  78. id := uuid.NewV4().String()
  79. mgr.st.addPending(id, e)
  80. go mgr.flushPending()
  81. }
  82. // leaderAccept assigns a sequence number and pushes the entry to every peer.
  83. func (mgr *Manager) leaderAccept(e Entry) Entry {
  84. lease := mgr.st.getLease()
  85. assigned := mgr.st.appendLog(e, lease.Term)
  86. mgr.st.setApplied(assigned.Seq, lease.Term)
  87. mgr.mu.Lock()
  88. mgr.appends++
  89. compact := mgr.appends%1000 == 0
  90. mgr.mu.Unlock()
  91. if compact {
  92. mgr.st.compactLog(mgr.opt.LogKeep)
  93. }
  94. go mgr.fanout([]Entry{assigned})
  95. return assigned
  96. }
  97. func (mgr *Manager) onlinePeers() []string {
  98. out := []string{}
  99. for _, n := range mgr.m.NodeViews() {
  100. if n.Local {
  101. continue
  102. }
  103. if n.State == membership.StateOnline || n.State == membership.StateDegraded {
  104. out = append(out, n.ID)
  105. }
  106. }
  107. return out
  108. }
  109. func (mgr *Manager) fanout(entries []Entry) {
  110. for _, peer := range mgr.onlinePeers() {
  111. go func(id string) {
  112. ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
  113. defer cancel()
  114. mgr.m.Transport().DoJSON(ctx, id, http.MethodPost, pathAppend, AppendRequest{Entries: entries}, nil)
  115. }(peer)
  116. }
  117. }
  118. // followerReceive handles entries pushed by the leader.
  119. func (mgr *Manager) followerReceive(entries []Entry) {
  120. applied, term := mgr.st.getApplied()
  121. for _, e := range entries {
  122. mgr.apply(e)
  123. if e.Term != term {
  124. //New leadership term: the numbering restarted, take a snapshot later
  125. term = e.Term
  126. applied = 0
  127. mgr.st.setApplied(0, term)
  128. mgr.requestCatchup()
  129. }
  130. if e.Seq == applied+1 {
  131. applied = e.Seq
  132. mgr.st.setApplied(applied, term)
  133. } else if e.Seq > applied+1 {
  134. mgr.requestCatchup()
  135. }
  136. mgr.st.setLastSeq(e.Seq)
  137. }
  138. }
  139. func (mgr *Manager) requestCatchup() {
  140. select {
  141. case mgr.catchupCh <- struct{}{}:
  142. default:
  143. }
  144. }
  145. func (mgr *Manager) catchupLoop() {
  146. defer mgr.wg.Done()
  147. for {
  148. select {
  149. case <-mgr.stop:
  150. return
  151. case <-mgr.catchupCh:
  152. mgr.catchup()
  153. }
  154. }
  155. }
  156. // catchup pulls missing log entries (or a snapshot) from the leader.
  157. func (mgr *Manager) catchup() {
  158. if mgr.IsLeader() {
  159. return
  160. }
  161. leader := mgr.Leader()
  162. if leader == "" {
  163. return
  164. }
  165. for i := 0; i < 100; i++ {
  166. applied, term := mgr.st.getApplied()
  167. if applied == 0 {
  168. //Fresh node or fresh term: the log may have been compacted below
  169. //what we need, a snapshot is the only safe starting point
  170. mgr.snapshotFromLeader(leader)
  171. return
  172. }
  173. ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
  174. resp, err := mgr.m.Transport().Do(ctx, leader, http.MethodGet, pathLog+"?after="+uitoa(applied), nil)
  175. cancel()
  176. if err != nil {
  177. return
  178. }
  179. if resp.Status == http.StatusGone {
  180. mgr.snapshotFromLeader(leader)
  181. return
  182. }
  183. if err := resp.Error(); err != nil {
  184. return
  185. }
  186. var page LogResponse
  187. if json.Unmarshal(resp.Body, &page) != nil {
  188. return
  189. }
  190. if page.Term != term {
  191. mgr.snapshotFromLeader(leader)
  192. return
  193. }
  194. for _, e := range page.Entries {
  195. mgr.apply(e)
  196. applied = e.Seq
  197. }
  198. mgr.st.setApplied(applied, page.Term)
  199. mgr.st.setLastSeq(page.LastSeq)
  200. if len(page.Entries) < logPageSize || applied >= page.LastSeq {
  201. return
  202. }
  203. }
  204. }
  205. func (mgr *Manager) snapshotFromLeader(leader string) {
  206. ctx, cancel := context.WithTimeout(context.Background(), 60*time.Second)
  207. defer cancel()
  208. var snap Snapshot
  209. if err := mgr.m.Transport().DoJSON(ctx, leader, http.MethodGet, pathSnapshot, nil, &snap); err != nil {
  210. return
  211. }
  212. mgr.applySnapshot(snap)
  213. }
  214. func (mgr *Manager) applySnapshot(snap Snapshot) {
  215. for i := range snap.Files {
  216. if mgr.st.putFile(&snap.Files[i]) && mgr.OnChange != nil {
  217. js, _ := json.Marshal(snap.Files[i])
  218. mgr.OnChange(KindFile, js)
  219. }
  220. }
  221. for i := range snap.Volumes {
  222. if mgr.st.putVolume(&snap.Volumes[i]) && mgr.OnChange != nil {
  223. js, _ := json.Marshal(snap.Volumes[i])
  224. mgr.OnChange(KindVolume, js)
  225. }
  226. }
  227. for i := range snap.Policies {
  228. if mgr.st.putPolicy(&snap.Policies[i]) && mgr.OnChange != nil {
  229. js, _ := json.Marshal(snap.Policies[i])
  230. mgr.OnChange(KindPolicy, js)
  231. }
  232. }
  233. for i := range snap.Jobs {
  234. if mgr.st.putJob(&snap.Jobs[i]) && mgr.OnChange != nil {
  235. js, _ := json.Marshal(snap.Jobs[i])
  236. mgr.OnChange(KindJob, js)
  237. }
  238. }
  239. for i := range snap.Settings {
  240. if mgr.st.putSetting(&snap.Settings[i]) && mgr.OnChange != nil {
  241. js, _ := json.Marshal(snap.Settings[i])
  242. mgr.OnChange(KindSetting, js)
  243. }
  244. }
  245. mgr.st.setApplied(snap.LastSeq, snap.Term)
  246. mgr.st.setLastSeq(snap.LastSeq)
  247. }
  248. // pendingLoop retries follower changes that no leader has accepted yet.
  249. func (mgr *Manager) pendingLoop() {
  250. defer mgr.wg.Done()
  251. ticker := time.NewTicker(mgr.opt.PendingRetry)
  252. defer ticker.Stop()
  253. for {
  254. select {
  255. case <-mgr.stop:
  256. return
  257. case <-ticker.C:
  258. mgr.flushPending()
  259. }
  260. }
  261. }
  262. // flushPending hands queued entries to the leader (or accepts them locally
  263. // when this node became the leader meanwhile).
  264. // flushPending drains the pending queue. It is called after every local
  265. // write, so only one call does the work at a time: a second call just asks
  266. // the running one for another round. Many concurrent drains each walking the
  267. // whole queue made the work grow with the square of the backlog.
  268. func (mgr *Manager) flushPending() {
  269. mgr.flushAgain.Store(true)
  270. for {
  271. if !mgr.flushing.CompareAndSwap(false, true) {
  272. return //the active drain will see flushAgain
  273. }
  274. for mgr.flushAgain.Swap(false) {
  275. mgr.drainPending()
  276. }
  277. mgr.flushing.Store(false)
  278. //A request that arrived as we were finishing gets its round now
  279. if !mgr.flushAgain.Load() {
  280. return
  281. }
  282. }
  283. }
  284. // drainPending sends every queued change to the leader, or accepts it when
  285. // this node is the leader.
  286. func (mgr *Manager) drainPending() {
  287. pending := mgr.st.pendingEntries()
  288. if len(pending) == 0 {
  289. return
  290. }
  291. if mgr.IsLeader() {
  292. for id, e := range pending {
  293. mgr.leaderAccept(e)
  294. mgr.st.removePending(id)
  295. }
  296. return
  297. }
  298. leader := mgr.Leader()
  299. if leader == "" {
  300. return
  301. }
  302. for id, e := range pending {
  303. ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second)
  304. err := mgr.m.Transport().DoJSON(ctx, leader, http.MethodPost, pathSubmit, e, nil)
  305. cancel()
  306. if err != nil {
  307. return
  308. }
  309. mgr.st.removePending(id)
  310. }
  311. }
  312. func uitoa(v uint64) string {
  313. if v == 0 {
  314. return "0"
  315. }
  316. buf := [20]byte{}
  317. i := len(buf)
  318. for v > 0 {
  319. i--
  320. buf[i] = byte('0' + v%10)
  321. v /= 10
  322. }
  323. return string(buf[i:])
  324. }