leveldb_backend.go 5.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233
  1. // Copyright (C) 2018 The Syncthing Authors.
  2. //
  3. // This Source Code Form is subject to the terms of the Mozilla Public
  4. // License, v. 2.0. If a copy of the MPL was not distributed with this file,
  5. // You can obtain one at https://mozilla.org/MPL/2.0/.
  6. package backend
  7. import (
  8. "github.com/syndtr/goleveldb/leveldb"
  9. "github.com/syndtr/goleveldb/leveldb/iterator"
  10. "github.com/syndtr/goleveldb/leveldb/util"
  11. )
  12. const (
  13. // Never flush transactions smaller than this, even on Checkpoint().
  14. // This just needs to be just large enough to avoid flushing
  15. // transactions when they are super tiny, thus creating millions of tiny
  16. // transactions unnecessarily.
  17. dbFlushBatchMin = 64 << KiB
  18. // Once a transaction reaches this size, flush it unconditionally. This
  19. // should be large enough to avoid forcing a flush between Checkpoint()
  20. // calls in loops where we do those, so in principle just large enough
  21. // to hold a FileInfo plus corresponding version list and metadata
  22. // updates or two.
  23. dbFlushBatchMax = 1 << MiB
  24. )
  25. // leveldbBackend implements Backend on top of a leveldb
  26. type leveldbBackend struct {
  27. ldb *leveldb.DB
  28. closeWG *closeWaitGroup
  29. location string
  30. }
  31. func newLeveldbBackend(ldb *leveldb.DB, location string) *leveldbBackend {
  32. return &leveldbBackend{
  33. ldb: ldb,
  34. closeWG: &closeWaitGroup{},
  35. location: location,
  36. }
  37. }
  38. func (b *leveldbBackend) NewReadTransaction() (ReadTransaction, error) {
  39. return b.newSnapshot()
  40. }
  41. func (b *leveldbBackend) newSnapshot() (leveldbSnapshot, error) {
  42. rel, err := newReleaser(b.closeWG)
  43. if err != nil {
  44. return leveldbSnapshot{}, err
  45. }
  46. snap, err := b.ldb.GetSnapshot()
  47. if err != nil {
  48. rel.Release()
  49. return leveldbSnapshot{}, wrapLeveldbErr(err)
  50. }
  51. return leveldbSnapshot{
  52. snap: snap,
  53. rel: rel,
  54. }, nil
  55. }
  56. func (b *leveldbBackend) NewWriteTransaction(hooks ...CommitHook) (WriteTransaction, error) {
  57. rel, err := newReleaser(b.closeWG)
  58. if err != nil {
  59. return nil, err
  60. }
  61. snap, err := b.newSnapshot()
  62. if err != nil {
  63. rel.Release()
  64. return nil, err // already wrapped
  65. }
  66. return &leveldbTransaction{
  67. leveldbSnapshot: snap,
  68. ldb: b.ldb,
  69. batch: new(leveldb.Batch),
  70. rel: rel,
  71. commitHooks: hooks,
  72. inFlush: false,
  73. }, nil
  74. }
  75. func (b *leveldbBackend) Close() error {
  76. b.closeWG.CloseWait()
  77. return wrapLeveldbErr(b.ldb.Close())
  78. }
  79. func (b *leveldbBackend) Get(key []byte) ([]byte, error) {
  80. val, err := b.ldb.Get(key, nil)
  81. return val, wrapLeveldbErr(err)
  82. }
  83. func (b *leveldbBackend) NewPrefixIterator(prefix []byte) (Iterator, error) {
  84. return &leveldbIterator{b.ldb.NewIterator(util.BytesPrefix(prefix), nil)}, nil
  85. }
  86. func (b *leveldbBackend) NewRangeIterator(first, last []byte) (Iterator, error) {
  87. return &leveldbIterator{b.ldb.NewIterator(&util.Range{Start: first, Limit: last}, nil)}, nil
  88. }
  89. func (b *leveldbBackend) Put(key, val []byte) error {
  90. return wrapLeveldbErr(b.ldb.Put(key, val, nil))
  91. }
  92. func (b *leveldbBackend) Delete(key []byte) error {
  93. return wrapLeveldbErr(b.ldb.Delete(key, nil))
  94. }
  95. func (b *leveldbBackend) Compact() error {
  96. // Race is detected during testing when db is closed while compaction
  97. // is ongoing.
  98. err := b.closeWG.Add(1)
  99. if err != nil {
  100. return err
  101. }
  102. defer b.closeWG.Done()
  103. return wrapLeveldbErr(b.ldb.CompactRange(util.Range{}))
  104. }
  105. func (b *leveldbBackend) Location() string {
  106. return b.location
  107. }
  108. // leveldbSnapshot implements backend.ReadTransaction
  109. type leveldbSnapshot struct {
  110. snap *leveldb.Snapshot
  111. rel *releaser
  112. }
  113. func (l leveldbSnapshot) Get(key []byte) ([]byte, error) {
  114. val, err := l.snap.Get(key, nil)
  115. return val, wrapLeveldbErr(err)
  116. }
  117. func (l leveldbSnapshot) NewPrefixIterator(prefix []byte) (Iterator, error) {
  118. return l.snap.NewIterator(util.BytesPrefix(prefix), nil), nil
  119. }
  120. func (l leveldbSnapshot) NewRangeIterator(first, last []byte) (Iterator, error) {
  121. return l.snap.NewIterator(&util.Range{Start: first, Limit: last}, nil), nil
  122. }
  123. func (l leveldbSnapshot) Release() {
  124. l.snap.Release()
  125. l.rel.Release()
  126. }
  127. // leveldbTransaction implements backend.WriteTransaction using a batch (not
  128. // an actual leveldb transaction)
  129. type leveldbTransaction struct {
  130. leveldbSnapshot
  131. ldb *leveldb.DB
  132. batch *leveldb.Batch
  133. rel *releaser
  134. commitHooks []CommitHook
  135. inFlush bool
  136. }
  137. func (t *leveldbTransaction) Delete(key []byte) error {
  138. t.batch.Delete(key)
  139. return t.checkFlush(dbFlushBatchMax)
  140. }
  141. func (t *leveldbTransaction) Put(key, val []byte) error {
  142. t.batch.Put(key, val)
  143. return t.checkFlush(dbFlushBatchMax)
  144. }
  145. func (t *leveldbTransaction) Checkpoint() error {
  146. return t.checkFlush(dbFlushBatchMin)
  147. }
  148. func (t *leveldbTransaction) Commit() error {
  149. err := wrapLeveldbErr(t.flush())
  150. t.leveldbSnapshot.Release()
  151. t.rel.Release()
  152. return err
  153. }
  154. func (t *leveldbTransaction) Release() {
  155. t.leveldbSnapshot.Release()
  156. t.rel.Release()
  157. }
  158. // checkFlush flushes and resets the batch if its size exceeds the given size.
  159. func (t *leveldbTransaction) checkFlush(size int) error {
  160. // Hooks might put values in the database, which triggers a checkFlush which might trigger a flush,
  161. // which might trigger the hooks.
  162. // Don't recurse...
  163. if t.inFlush || len(t.batch.Dump()) < size {
  164. return nil
  165. }
  166. return t.flush()
  167. }
  168. func (t *leveldbTransaction) flush() error {
  169. t.inFlush = true
  170. defer func() { t.inFlush = false }()
  171. for _, hook := range t.commitHooks {
  172. if err := hook(t); err != nil {
  173. return err
  174. }
  175. }
  176. if t.batch.Len() == 0 {
  177. return nil
  178. }
  179. if err := t.ldb.Write(t.batch, nil); err != nil {
  180. return wrapLeveldbErr(err)
  181. }
  182. t.batch.Reset()
  183. return nil
  184. }
  185. type leveldbIterator struct {
  186. iterator.Iterator
  187. }
  188. func (it *leveldbIterator) Error() error {
  189. return wrapLeveldbErr(it.Iterator.Error())
  190. }
  191. // wrapLeveldbErr wraps errors so that the backend package can recognize them
  192. func wrapLeveldbErr(err error) error {
  193. if err == leveldb.ErrClosed {
  194. return &errClosed{}
  195. }
  196. if err == leveldb.ErrNotFound {
  197. return &errNotFound{}
  198. }
  199. return err
  200. }