db_bolt.go 4.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192
  1. package manager
  2. import (
  3. "encoding/json"
  4. "fmt"
  5. "strings"
  6. "time"
  7. "github.com/boltdb/bolt"
  8. . "github.com/tuna/tunasync/internal"
  9. )
  10. const (
  11. _workerBucketKey = "workers"
  12. _statusBucketKey = "mirror_status"
  13. )
  14. type boltAdapter struct {
  15. db *bolt.DB
  16. dbFile string
  17. }
  18. func (b *boltAdapter) Init() (err error) {
  19. return b.db.Update(func(tx *bolt.Tx) error {
  20. _, err = tx.CreateBucketIfNotExists([]byte(_workerBucketKey))
  21. if err != nil {
  22. return fmt.Errorf("create bucket %s error: %s", _workerBucketKey, err.Error())
  23. }
  24. _, err = tx.CreateBucketIfNotExists([]byte(_statusBucketKey))
  25. if err != nil {
  26. return fmt.Errorf("create bucket %s error: %s", _statusBucketKey, err.Error())
  27. }
  28. return nil
  29. })
  30. }
  31. func (b *boltAdapter) ListWorkers() (ws []WorkerStatus, err error) {
  32. err = b.db.View(func(tx *bolt.Tx) error {
  33. bucket := tx.Bucket([]byte(_workerBucketKey))
  34. c := bucket.Cursor()
  35. var w WorkerStatus
  36. for k, v := c.First(); k != nil; k, v = c.Next() {
  37. jsonErr := json.Unmarshal(v, &w)
  38. if jsonErr != nil {
  39. err = fmt.Errorf("%s; %s", err.Error(), jsonErr)
  40. continue
  41. }
  42. ws = append(ws, w)
  43. }
  44. return err
  45. })
  46. return
  47. }
  48. func (b *boltAdapter) GetWorker(workerID string) (w WorkerStatus, err error) {
  49. err = b.db.View(func(tx *bolt.Tx) error {
  50. bucket := tx.Bucket([]byte(_workerBucketKey))
  51. v := bucket.Get([]byte(workerID))
  52. if v == nil {
  53. return fmt.Errorf("invalid workerID %s", workerID)
  54. }
  55. err := json.Unmarshal(v, &w)
  56. return err
  57. })
  58. return
  59. }
  60. func (b *boltAdapter) DeleteWorker(workerID string) (err error) {
  61. err = b.db.Update(func(tx *bolt.Tx) error {
  62. bucket := tx.Bucket([]byte(_workerBucketKey))
  63. v := bucket.Get([]byte(workerID))
  64. if v == nil {
  65. return fmt.Errorf("invalid workerID %s", workerID)
  66. }
  67. err := bucket.Delete([]byte(workerID))
  68. return err
  69. })
  70. return
  71. }
  72. func (b *boltAdapter) CreateWorker(w WorkerStatus) (WorkerStatus, error) {
  73. err := b.db.Update(func(tx *bolt.Tx) error {
  74. bucket := tx.Bucket([]byte(_workerBucketKey))
  75. v, err := json.Marshal(w)
  76. if err != nil {
  77. return err
  78. }
  79. err = bucket.Put([]byte(w.ID), v)
  80. return err
  81. })
  82. return w, err
  83. }
  84. func (b *boltAdapter) RefreshWorker(workerID string) (w WorkerStatus, err error) {
  85. w, err = b.GetWorker(workerID)
  86. if err == nil {
  87. w.LastOnline = time.Now()
  88. w, err = b.CreateWorker(w)
  89. }
  90. return w, err
  91. }
  92. func (b *boltAdapter) UpdateMirrorStatus(workerID, mirrorID string, status MirrorStatus) (MirrorStatus, error) {
  93. id := mirrorID + "/" + workerID
  94. err := b.db.Update(func(tx *bolt.Tx) error {
  95. bucket := tx.Bucket([]byte(_statusBucketKey))
  96. v, err := json.Marshal(status)
  97. err = bucket.Put([]byte(id), v)
  98. return err
  99. })
  100. return status, err
  101. }
  102. func (b *boltAdapter) GetMirrorStatus(workerID, mirrorID string) (m MirrorStatus, err error) {
  103. id := mirrorID + "/" + workerID
  104. err = b.db.Update(func(tx *bolt.Tx) error {
  105. bucket := tx.Bucket([]byte(_statusBucketKey))
  106. v := bucket.Get([]byte(id))
  107. if v == nil {
  108. return fmt.Errorf("no mirror '%s' exists in worker '%s'", mirrorID, workerID)
  109. }
  110. err := json.Unmarshal(v, &m)
  111. return err
  112. })
  113. return
  114. }
  115. func (b *boltAdapter) ListMirrorStatus(workerID string) (ms []MirrorStatus, err error) {
  116. err = b.db.View(func(tx *bolt.Tx) error {
  117. bucket := tx.Bucket([]byte(_statusBucketKey))
  118. c := bucket.Cursor()
  119. var m MirrorStatus
  120. for k, v := c.First(); k != nil; k, v = c.Next() {
  121. if wID := strings.Split(string(k), "/")[1]; wID == workerID {
  122. jsonErr := json.Unmarshal(v, &m)
  123. if jsonErr != nil {
  124. err = fmt.Errorf("%s; %s", err.Error(), jsonErr)
  125. continue
  126. }
  127. ms = append(ms, m)
  128. }
  129. }
  130. return err
  131. })
  132. return
  133. }
  134. func (b *boltAdapter) ListAllMirrorStatus() (ms []MirrorStatus, err error) {
  135. err = b.db.View(func(tx *bolt.Tx) error {
  136. bucket := tx.Bucket([]byte(_statusBucketKey))
  137. c := bucket.Cursor()
  138. var m MirrorStatus
  139. for k, v := c.First(); k != nil; k, v = c.Next() {
  140. jsonErr := json.Unmarshal(v, &m)
  141. if jsonErr != nil {
  142. err = fmt.Errorf("%s; %s", err.Error(), jsonErr)
  143. continue
  144. }
  145. ms = append(ms, m)
  146. }
  147. return err
  148. })
  149. return
  150. }
  151. func (b *boltAdapter) FlushDisabledJobs() (err error) {
  152. err = b.db.Update(func(tx *bolt.Tx) error {
  153. bucket := tx.Bucket([]byte(_statusBucketKey))
  154. c := bucket.Cursor()
  155. var m MirrorStatus
  156. for k, v := c.First(); k != nil; k, v = c.Next() {
  157. jsonErr := json.Unmarshal(v, &m)
  158. if jsonErr != nil {
  159. err = fmt.Errorf("%s; %s", err.Error(), jsonErr)
  160. continue
  161. }
  162. if m.Status == Disabled || len(m.Name) == 0 {
  163. err = c.Delete()
  164. }
  165. }
  166. return err
  167. })
  168. return
  169. }
  170. func (b *boltAdapter) Close() error {
  171. if b.db != nil {
  172. return b.db.Close()
  173. }
  174. return nil
  175. }