lib/db: Add closeWaitGroup to allow async operation (#6317)
This commit is contained in:
@@ -150,16 +150,18 @@ func IsNotFound(err error) bool {
|
||||
|
||||
// releaser manages counting on top of a waitgroup
|
||||
type releaser struct {
|
||||
wg *sync.WaitGroup
|
||||
wg *closeWaitGroup
|
||||
once *sync.Once
|
||||
}
|
||||
|
||||
func newReleaser(wg *sync.WaitGroup) *releaser {
|
||||
wg.Add(1)
|
||||
func newReleaser(wg *closeWaitGroup) (*releaser, error) {
|
||||
if err := wg.Add(1); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &releaser{
|
||||
wg: wg,
|
||||
once: new(sync.Once),
|
||||
}
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (r releaser) Release() {
|
||||
@@ -169,3 +171,29 @@ func (r releaser) Release() {
|
||||
r.wg.Done()
|
||||
})
|
||||
}
|
||||
|
||||
// closeWaitGroup behaves just like a sync.WaitGroup, but does not require
|
||||
// a single routine to do the Add and Wait calls. If Add is called after
|
||||
// CloseWait, it will return an error, and both are safe to be used concurrently.
|
||||
type closeWaitGroup struct {
|
||||
sync.WaitGroup
|
||||
closed bool
|
||||
closeMut sync.RWMutex
|
||||
}
|
||||
|
||||
func (cg *closeWaitGroup) Add(i int) error {
|
||||
cg.closeMut.RLock()
|
||||
defer cg.closeMut.RUnlock()
|
||||
if cg.closed {
|
||||
return errClosed{}
|
||||
}
|
||||
cg.WaitGroup.Add(i)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (cg *closeWaitGroup) CloseWait() {
|
||||
cg.closeMut.Lock()
|
||||
cg.closed = true
|
||||
cg.closeMut.Unlock()
|
||||
cg.WaitGroup.Wait()
|
||||
}
|
||||
|
||||
@@ -7,8 +7,6 @@
|
||||
package backend
|
||||
|
||||
import (
|
||||
"sync"
|
||||
|
||||
"github.com/syndtr/goleveldb/leveldb"
|
||||
"github.com/syndtr/goleveldb/leveldb/iterator"
|
||||
"github.com/syndtr/goleveldb/leveldb/util"
|
||||
@@ -24,7 +22,14 @@ const (
|
||||
// leveldbBackend implements Backend on top of a leveldb
|
||||
type leveldbBackend struct {
|
||||
ldb *leveldb.DB
|
||||
closeWG sync.WaitGroup
|
||||
closeWG *closeWaitGroup
|
||||
}
|
||||
|
||||
func newLeveldbBackend(ldb *leveldb.DB) *leveldbBackend {
|
||||
return &leveldbBackend{
|
||||
ldb: ldb,
|
||||
closeWG: &closeWaitGroup{},
|
||||
}
|
||||
}
|
||||
|
||||
func (b *leveldbBackend) NewReadTransaction() (ReadTransaction, error) {
|
||||
@@ -36,9 +41,13 @@ func (b *leveldbBackend) newSnapshot() (leveldbSnapshot, error) {
|
||||
if err != nil {
|
||||
return leveldbSnapshot{}, wrapLeveldbErr(err)
|
||||
}
|
||||
rel, err := newReleaser(b.closeWG)
|
||||
if err != nil {
|
||||
return leveldbSnapshot{}, err
|
||||
}
|
||||
return leveldbSnapshot{
|
||||
snap: snap,
|
||||
rel: newReleaser(&b.closeWG),
|
||||
rel: rel,
|
||||
}, nil
|
||||
}
|
||||
|
||||
@@ -47,16 +56,20 @@ func (b *leveldbBackend) NewWriteTransaction() (WriteTransaction, error) {
|
||||
if err != nil {
|
||||
return nil, err // already wrapped
|
||||
}
|
||||
rel, err := newReleaser(b.closeWG)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &leveldbTransaction{
|
||||
leveldbSnapshot: snap,
|
||||
ldb: b.ldb,
|
||||
batch: new(leveldb.Batch),
|
||||
rel: newReleaser(&b.closeWG),
|
||||
rel: rel,
|
||||
}, nil
|
||||
}
|
||||
|
||||
func (b *leveldbBackend) Close() error {
|
||||
b.closeWG.Wait()
|
||||
b.closeWG.CloseWait()
|
||||
return wrapLeveldbErr(b.ldb.Close())
|
||||
}
|
||||
|
||||
@@ -82,6 +95,13 @@ func (b *leveldbBackend) Delete(key []byte) error {
|
||||
}
|
||||
|
||||
func (b *leveldbBackend) Compact() error {
|
||||
// Race is detected during testing when db is closed while compaction
|
||||
// is ongoing.
|
||||
err := b.closeWG.Add(1)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
defer b.closeWG.Done()
|
||||
return wrapLeveldbErr(b.ldb.CompactRange(util.Range{}))
|
||||
}
|
||||
|
||||
|
||||
@@ -42,7 +42,7 @@ func OpenLevelDB(location string, tuning Tuning) (Backend, error) {
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &leveldbBackend{ldb: ldb}, nil
|
||||
return newLeveldbBackend(ldb), nil
|
||||
}
|
||||
|
||||
// OpenRO attempts to open the database at the given location, read only.
|
||||
@@ -55,13 +55,13 @@ func OpenLevelDBRO(location string) (Backend, error) {
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return &leveldbBackend{ldb: ldb}, nil
|
||||
return newLeveldbBackend(ldb), nil
|
||||
}
|
||||
|
||||
// OpenMemory returns a new Backend referencing an in-memory database.
|
||||
func OpenLevelDBMemory() Backend {
|
||||
ldb, _ := leveldb.Open(storage.NewMemStorage(), nil)
|
||||
return &leveldbBackend{ldb: ldb}
|
||||
return newLeveldbBackend(ldb)
|
||||
}
|
||||
|
||||
// optsFor returns the database options to use when opening a database with
|
||||
|
||||
Reference in New Issue
Block a user