db.go 5.3 KB

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