leveldb_backend.go 4.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173
  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. "sync"
  9. "github.com/syndtr/goleveldb/leveldb"
  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 sync.WaitGroup
  22. }
  23. func (b *leveldbBackend) NewReadTransaction() (ReadTransaction, error) {
  24. return b.newSnapshot()
  25. }
  26. func (b *leveldbBackend) newSnapshot() (leveldbSnapshot, error) {
  27. snap, err := b.ldb.GetSnapshot()
  28. if err != nil {
  29. return leveldbSnapshot{}, wrapLeveldbErr(err)
  30. }
  31. return leveldbSnapshot{
  32. snap: snap,
  33. rel: newReleaser(&b.closeWG),
  34. }, nil
  35. }
  36. func (b *leveldbBackend) NewWriteTransaction() (WriteTransaction, error) {
  37. snap, err := b.newSnapshot()
  38. if err != nil {
  39. return nil, err // already wrapped
  40. }
  41. return &leveldbTransaction{
  42. leveldbSnapshot: snap,
  43. ldb: b.ldb,
  44. batch: new(leveldb.Batch),
  45. rel: newReleaser(&b.closeWG),
  46. }, nil
  47. }
  48. func (b *leveldbBackend) Close() error {
  49. b.closeWG.Wait()
  50. return wrapLeveldbErr(b.ldb.Close())
  51. }
  52. func (b *leveldbBackend) Get(key []byte) ([]byte, error) {
  53. val, err := b.ldb.Get(key, nil)
  54. return val, wrapLeveldbErr(err)
  55. }
  56. func (b *leveldbBackend) NewPrefixIterator(prefix []byte) (Iterator, error) {
  57. return b.ldb.NewIterator(util.BytesPrefix(prefix), nil), nil
  58. }
  59. func (b *leveldbBackend) NewRangeIterator(first, last []byte) (Iterator, error) {
  60. return b.ldb.NewIterator(&util.Range{Start: first, Limit: last}, nil), nil
  61. }
  62. func (b *leveldbBackend) Put(key, val []byte) error {
  63. return wrapLeveldbErr(b.ldb.Put(key, val, nil))
  64. }
  65. func (b *leveldbBackend) Delete(key []byte) error {
  66. return wrapLeveldbErr(b.ldb.Delete(key, nil))
  67. }
  68. // leveldbSnapshot implements backend.ReadTransaction
  69. type leveldbSnapshot struct {
  70. snap *leveldb.Snapshot
  71. rel *releaser
  72. }
  73. func (l leveldbSnapshot) Get(key []byte) ([]byte, error) {
  74. val, err := l.snap.Get(key, nil)
  75. return val, wrapLeveldbErr(err)
  76. }
  77. func (l leveldbSnapshot) NewPrefixIterator(prefix []byte) (Iterator, error) {
  78. return l.snap.NewIterator(util.BytesPrefix(prefix), nil), nil
  79. }
  80. func (l leveldbSnapshot) NewRangeIterator(first, last []byte) (Iterator, error) {
  81. return l.snap.NewIterator(&util.Range{Start: first, Limit: last}, nil), nil
  82. }
  83. func (l leveldbSnapshot) Release() {
  84. l.snap.Release()
  85. l.rel.Release()
  86. }
  87. // leveldbTransaction implements backend.WriteTransaction using a batch (not
  88. // an actual leveldb transaction)
  89. type leveldbTransaction struct {
  90. leveldbSnapshot
  91. ldb *leveldb.DB
  92. batch *leveldb.Batch
  93. rel *releaser
  94. }
  95. func (t *leveldbTransaction) Delete(key []byte) error {
  96. t.batch.Delete(key)
  97. return t.checkFlush(dbFlushBatchMax)
  98. }
  99. func (t *leveldbTransaction) Put(key, val []byte) error {
  100. t.batch.Put(key, val)
  101. return t.checkFlush(dbFlushBatchMax)
  102. }
  103. func (t *leveldbTransaction) Checkpoint() error {
  104. return t.checkFlush(dbFlushBatchMin)
  105. }
  106. func (t *leveldbTransaction) Commit() error {
  107. err := wrapLeveldbErr(t.flush())
  108. t.leveldbSnapshot.Release()
  109. t.rel.Release()
  110. return err
  111. }
  112. func (t *leveldbTransaction) Release() {
  113. t.leveldbSnapshot.Release()
  114. t.rel.Release()
  115. }
  116. // checkFlush flushes and resets the batch if its size exceeds the given size.
  117. func (t *leveldbTransaction) checkFlush(size int) error {
  118. if len(t.batch.Dump()) < size {
  119. return nil
  120. }
  121. return t.flush()
  122. }
  123. func (t *leveldbTransaction) flush() error {
  124. if t.batch.Len() == 0 {
  125. return nil
  126. }
  127. if err := t.ldb.Write(t.batch, nil); err != nil {
  128. return wrapLeveldbErr(err)
  129. }
  130. t.batch.Reset()
  131. return nil
  132. }
  133. // wrapLeveldbErr wraps errors so that the backend package can recognize them
  134. func wrapLeveldbErr(err error) error {
  135. if err == nil {
  136. return nil
  137. }
  138. if err == leveldb.ErrClosed {
  139. return errClosed{}
  140. }
  141. if err == leveldb.ErrNotFound {
  142. return errNotFound{}
  143. }
  144. return err
  145. }