leveldb_backend.go 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213
  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. dbFlushBatchMin = 1 << MiB
  15. // Once a transaction reaches this size, flush it unconditionally.
  16. dbFlushBatchMax = 128 << MiB
  17. )
  18. // leveldbBackend implements Backend on top of a leveldb
  19. type leveldbBackend struct {
  20. ldb *leveldb.DB
  21. closeWG *closeWaitGroup
  22. }
  23. func newLeveldbBackend(ldb *leveldb.DB) *leveldbBackend {
  24. return &leveldbBackend{
  25. ldb: ldb,
  26. closeWG: &closeWaitGroup{},
  27. }
  28. }
  29. func (b *leveldbBackend) NewReadTransaction() (ReadTransaction, error) {
  30. return b.newSnapshot()
  31. }
  32. func (b *leveldbBackend) newSnapshot() (leveldbSnapshot, error) {
  33. rel, err := newReleaser(b.closeWG)
  34. if err != nil {
  35. return leveldbSnapshot{}, err
  36. }
  37. snap, err := b.ldb.GetSnapshot()
  38. if err != nil {
  39. rel.Release()
  40. return leveldbSnapshot{}, wrapLeveldbErr(err)
  41. }
  42. return leveldbSnapshot{
  43. snap: snap,
  44. rel: rel,
  45. }, nil
  46. }
  47. func (b *leveldbBackend) NewWriteTransaction() (WriteTransaction, error) {
  48. rel, err := newReleaser(b.closeWG)
  49. if err != nil {
  50. return nil, err
  51. }
  52. snap, err := b.newSnapshot()
  53. if err != nil {
  54. rel.Release()
  55. return nil, err // already wrapped
  56. }
  57. return &leveldbTransaction{
  58. leveldbSnapshot: snap,
  59. ldb: b.ldb,
  60. batch: new(leveldb.Batch),
  61. rel: rel,
  62. }, nil
  63. }
  64. func (b *leveldbBackend) Close() error {
  65. b.closeWG.CloseWait()
  66. return wrapLeveldbErr(b.ldb.Close())
  67. }
  68. func (b *leveldbBackend) Get(key []byte) ([]byte, error) {
  69. val, err := b.ldb.Get(key, nil)
  70. return val, wrapLeveldbErr(err)
  71. }
  72. func (b *leveldbBackend) NewPrefixIterator(prefix []byte) (Iterator, error) {
  73. return &leveldbIterator{b.ldb.NewIterator(util.BytesPrefix(prefix), nil)}, nil
  74. }
  75. func (b *leveldbBackend) NewRangeIterator(first, last []byte) (Iterator, error) {
  76. return &leveldbIterator{b.ldb.NewIterator(&util.Range{Start: first, Limit: last}, nil)}, nil
  77. }
  78. func (b *leveldbBackend) Put(key, val []byte) error {
  79. return wrapLeveldbErr(b.ldb.Put(key, val, nil))
  80. }
  81. func (b *leveldbBackend) Delete(key []byte) error {
  82. return wrapLeveldbErr(b.ldb.Delete(key, nil))
  83. }
  84. func (b *leveldbBackend) Compact() error {
  85. // Race is detected during testing when db is closed while compaction
  86. // is ongoing.
  87. err := b.closeWG.Add(1)
  88. if err != nil {
  89. return err
  90. }
  91. defer b.closeWG.Done()
  92. return wrapLeveldbErr(b.ldb.CompactRange(util.Range{}))
  93. }
  94. // leveldbSnapshot implements backend.ReadTransaction
  95. type leveldbSnapshot struct {
  96. snap *leveldb.Snapshot
  97. rel *releaser
  98. }
  99. func (l leveldbSnapshot) Get(key []byte) ([]byte, error) {
  100. val, err := l.snap.Get(key, nil)
  101. return val, wrapLeveldbErr(err)
  102. }
  103. func (l leveldbSnapshot) NewPrefixIterator(prefix []byte) (Iterator, error) {
  104. return l.snap.NewIterator(util.BytesPrefix(prefix), nil), nil
  105. }
  106. func (l leveldbSnapshot) NewRangeIterator(first, last []byte) (Iterator, error) {
  107. return l.snap.NewIterator(&util.Range{Start: first, Limit: last}, nil), nil
  108. }
  109. func (l leveldbSnapshot) Release() {
  110. l.snap.Release()
  111. l.rel.Release()
  112. }
  113. // leveldbTransaction implements backend.WriteTransaction using a batch (not
  114. // an actual leveldb transaction)
  115. type leveldbTransaction struct {
  116. leveldbSnapshot
  117. ldb *leveldb.DB
  118. batch *leveldb.Batch
  119. rel *releaser
  120. }
  121. func (t *leveldbTransaction) Delete(key []byte) error {
  122. t.batch.Delete(key)
  123. return t.checkFlush(dbFlushBatchMax)
  124. }
  125. func (t *leveldbTransaction) Put(key, val []byte) error {
  126. t.batch.Put(key, val)
  127. return t.checkFlush(dbFlushBatchMax)
  128. }
  129. func (t *leveldbTransaction) Checkpoint(preFlush ...func() error) error {
  130. return t.checkFlush(dbFlushBatchMin, preFlush...)
  131. }
  132. func (t *leveldbTransaction) Commit() error {
  133. err := wrapLeveldbErr(t.flush())
  134. t.leveldbSnapshot.Release()
  135. t.rel.Release()
  136. return err
  137. }
  138. func (t *leveldbTransaction) Release() {
  139. t.leveldbSnapshot.Release()
  140. t.rel.Release()
  141. }
  142. // checkFlush flushes and resets the batch if its size exceeds the given size.
  143. func (t *leveldbTransaction) checkFlush(size int, preFlush ...func() error) error {
  144. if len(t.batch.Dump()) < size {
  145. return nil
  146. }
  147. for _, hook := range preFlush {
  148. if err := hook(); err != nil {
  149. return err
  150. }
  151. }
  152. return t.flush()
  153. }
  154. func (t *leveldbTransaction) flush() error {
  155. if t.batch.Len() == 0 {
  156. return nil
  157. }
  158. if err := t.ldb.Write(t.batch, nil); err != nil {
  159. return wrapLeveldbErr(err)
  160. }
  161. t.batch.Reset()
  162. return nil
  163. }
  164. type leveldbIterator struct {
  165. iterator.Iterator
  166. }
  167. func (it *leveldbIterator) Error() error {
  168. return wrapLeveldbErr(it.Iterator.Error())
  169. }
  170. // wrapLeveldbErr wraps errors so that the backend package can recognize them
  171. func wrapLeveldbErr(err error) error {
  172. if err == nil {
  173. return nil
  174. }
  175. if err == leveldb.ErrClosed {
  176. return errClosed{}
  177. }
  178. if err == leveldb.ErrNotFound {
  179. return errNotFound{}
  180. }
  181. return err
  182. }