persist.go 4.5 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184
  1. package metadata
  2. /*
  3. Write-behind persistence for the metadata store.
  4. The store keeps everything it serves in memory; the disk copy exists so
  5. a node can restart without a full resync. Writing that copy inline, one
  6. synced transaction per change while holding the store lock, made every
  7. reader and writer queue behind the disk. Under a burst the lease renewal
  8. waited long enough for the leader to lose its lease, and the node never
  9. caught up.
  10. So a change now updates memory and queues its disk write here. Queued
  11. writes are coalesced per key (only the newest value of a key is written)
  12. and a background writer commits the whole queue in one transaction every
  13. FlushInterval. Close flushes what is left.
  14. A crash can lose the last FlushInterval of changes from this node's disk
  15. copy. That copy is not authoritative: the replicated log and snapshots
  16. bring a restarted node back up to date, exactly as they do for a node
  17. that was offline.
  18. */
  19. import (
  20. "encoding/json"
  21. "sync"
  22. "time"
  23. "imuslab.com/arozos/mod/database"
  24. "imuslab.com/arozos/mod/info/logger"
  25. )
  26. // FlushInterval is how often queued writes reach the disk.
  27. var FlushInterval = 50 * time.Millisecond
  28. type opKey struct{ table, key string }
  29. type persister struct {
  30. db *database.Database
  31. interval time.Duration //FlushInterval when the persister was made
  32. mu sync.Mutex
  33. queue map[opKey]database.BatchOp
  34. stopped bool
  35. flushMu sync.Mutex //one commit at a time, so older values never land last
  36. once sync.Once //close or discard, whichever comes first
  37. wake chan struct{}
  38. stop chan struct{}
  39. done chan struct{}
  40. }
  41. func newPersister(db *database.Database) *persister {
  42. p := &persister{
  43. db: db,
  44. interval: FlushInterval,
  45. queue: map[opKey]database.BatchOp{},
  46. wake: make(chan struct{}, 1),
  47. stop: make(chan struct{}),
  48. done: make(chan struct{}),
  49. }
  50. go p.loop()
  51. return p
  52. }
  53. // put queues a write. The value is marshalled now, so the caller may keep
  54. // changing its own copy after put returns.
  55. func (p *persister) put(table string, key string, value interface{}) {
  56. js, err := json.Marshal(value)
  57. if err != nil {
  58. logger.PrintAndLog("Cluster", "Metadata "+table+"/"+key+" could not be encoded", err)
  59. return
  60. }
  61. p.enqueue(database.BatchOp{Table: table, Key: key, Value: json.RawMessage(js)})
  62. }
  63. // del queues a removal.
  64. func (p *persister) del(table string, key string) {
  65. p.enqueue(database.BatchOp{Table: table, Key: key, Delete: true})
  66. }
  67. func (p *persister) enqueue(op database.BatchOp) {
  68. p.mu.Lock()
  69. if p.stopped {
  70. p.mu.Unlock()
  71. return
  72. }
  73. p.queue[opKey{op.Table, op.Key}] = op //newest change of a key wins
  74. p.mu.Unlock()
  75. select {
  76. case p.wake <- struct{}{}:
  77. default:
  78. }
  79. }
  80. func (p *persister) loop() {
  81. defer close(p.done)
  82. for {
  83. select {
  84. case <-p.stop:
  85. return
  86. case <-p.wake:
  87. }
  88. //Let a burst gather so it shares one transaction
  89. select {
  90. case <-p.stop:
  91. return
  92. case <-time.After(p.interval):
  93. }
  94. p.flush()
  95. }
  96. }
  97. // flush commits everything queued so far and returns once it is on disk.
  98. func (p *persister) flush() {
  99. p.flushMu.Lock()
  100. defer p.flushMu.Unlock()
  101. p.mu.Lock()
  102. if len(p.queue) == 0 {
  103. p.mu.Unlock()
  104. return
  105. }
  106. batch := p.queue
  107. p.queue = map[opKey]database.BatchOp{}
  108. p.mu.Unlock()
  109. ops := make([]database.BatchOp, 0, len(batch))
  110. for _, op := range batch {
  111. ops = append(ops, op)
  112. }
  113. if err := p.db.WriteBatch(ops); err != nil {
  114. logger.PrintAndLog("Cluster", "Metadata could not be written to disk; retrying", err)
  115. //Put the batch back unless a newer change of the same key arrived
  116. p.mu.Lock()
  117. if !p.stopped {
  118. for k, op := range batch {
  119. if _, newer := p.queue[k]; !newer {
  120. p.queue[k] = op
  121. }
  122. }
  123. }
  124. p.mu.Unlock()
  125. select {
  126. case p.wake <- struct{}{}:
  127. default:
  128. }
  129. }
  130. }
  131. // pending reports how many keys wait to be written (for tests and status).
  132. func (p *persister) pending() int {
  133. p.mu.Lock()
  134. defer p.mu.Unlock()
  135. return len(p.queue)
  136. }
  137. // close stops the writer and commits what is left.
  138. func (p *persister) close() {
  139. p.once.Do(func() {
  140. close(p.stop)
  141. <-p.done
  142. p.mu.Lock()
  143. p.stopped = true
  144. p.mu.Unlock()
  145. p.flush()
  146. })
  147. }
  148. // discard stops the writer and drops what is queued, for a store whose
  149. // tables were wiped: writing its leftovers would bring the old data back.
  150. func (p *persister) discard() {
  151. p.once.Do(func() {
  152. p.mu.Lock()
  153. p.stopped = true
  154. p.queue = map[opKey]database.BatchOp{}
  155. p.mu.Unlock()
  156. close(p.stop)
  157. <-p.done
  158. //Wait out a commit another caller had already started
  159. p.flushMu.Lock()
  160. p.flushMu.Unlock()
  161. })
  162. }