swamr/storage/localstore: wait for gc and writeGCSize workers on Close

This commit is contained in:
Janos Guljas 2019-02-21 17:50:27 +01:00
parent 9f5e9d6d94
commit f25b2e026e
2 changed files with 30 additions and 3 deletions

View file

@ -117,6 +117,8 @@ var (
// run. GC run iterates on gcIndex and removes older items // run. GC run iterates on gcIndex and removes older items
// form retrieval and other indexes. // form retrieval and other indexes.
func (db *DB) collectGarbageWorker() { func (db *DB) collectGarbageWorker() {
defer close(db.collectGarbageWorkerDone)
for { for {
select { select {
case <-db.collectGarbageTrigger: case <-db.collectGarbageTrigger:
@ -243,6 +245,8 @@ func (db *DB) triggerGarbageCollection() {
// writeGCSizeDelay duration to avoid very frequent // writeGCSizeDelay duration to avoid very frequent
// database operations. // database operations.
func (db *DB) writeGCSizeWorker() { func (db *DB) writeGCSizeWorker() {
defer close(db.writeGCSizeWorkerDone)
for { for {
select { select {
case <-db.writeGCSizeTrigger: case <-db.writeGCSizeTrigger:

View file

@ -107,6 +107,12 @@ type DB struct {
// this channel is closed when close function is called // this channel is closed when close function is called
// to terminate other goroutines // to terminate other goroutines
close chan struct{} close chan struct{}
// protect Close method from exiting before
// garbage collection and gc size write workers
// are done
collectGarbageWorkerDone chan struct{}
writeGCSizeWorkerDone chan struct{}
} }
// Options struct holds optional parameters for configuring DB. // Options struct holds optional parameters for configuring DB.
@ -138,9 +144,11 @@ func New(path string, baseKey []byte, o *Options) (db *DB, err error) {
// need to be buffered with the size of 1 // need to be buffered with the size of 1
// to signal another event if it // to signal another event if it
// is triggered during already running function // is triggered during already running function
collectGarbageTrigger: make(chan struct{}, 1), collectGarbageTrigger: make(chan struct{}, 1),
writeGCSizeTrigger: make(chan struct{}, 1), writeGCSizeTrigger: make(chan struct{}, 1),
close: make(chan struct{}), close: make(chan struct{}),
collectGarbageWorkerDone: make(chan struct{}),
writeGCSizeWorkerDone: make(chan struct{}),
} }
if db.capacity <= 0 { if db.capacity <= 0 {
db.capacity = defaultCapacity db.capacity = defaultCapacity
@ -361,6 +369,21 @@ func New(path string, baseKey []byte, o *Options) (db *DB, err error) {
func (db *DB) Close() (err error) { func (db *DB) Close() (err error) {
close(db.close) close(db.close)
db.updateGCWG.Wait() db.updateGCWG.Wait()
// wait for gc worker and gc size write workers to
// return before closing the shed
timeout := time.After(10 * time.Second)
select {
case <-db.collectGarbageWorkerDone:
case <-timeout:
log.Error("localstore: collect garbage worker did not return after db close")
}
select {
case <-db.writeGCSizeWorkerDone:
case <-timeout:
log.Error("localstore: write gc size worker did not return after db close")
}
if err := db.writeGCSize(db.getGCSize()); err != nil { if err := db.writeGCSize(db.getGCSize()); err != nil {
log.Error("localstore: write gc size", "err", err) log.Error("localstore: write gc size", "err", err)
} }