From f25b2e026ea08d91d57fc8c79b04ea4b0e8cef5e Mon Sep 17 00:00:00 2001 From: Janos Guljas Date: Thu, 21 Feb 2019 17:50:27 +0100 Subject: [PATCH] swamr/storage/localstore: wait for gc and writeGCSize workers on Close --- swarm/storage/localstore/gc.go | 4 ++++ swarm/storage/localstore/localstore.go | 29 +++++++++++++++++++++++--- 2 files changed, 30 insertions(+), 3 deletions(-) diff --git a/swarm/storage/localstore/gc.go b/swarm/storage/localstore/gc.go index 7718d1e589..50fbfe6fec 100644 --- a/swarm/storage/localstore/gc.go +++ b/swarm/storage/localstore/gc.go @@ -117,6 +117,8 @@ var ( // run. GC run iterates on gcIndex and removes older items // form retrieval and other indexes. func (db *DB) collectGarbageWorker() { + defer close(db.collectGarbageWorkerDone) + for { select { case <-db.collectGarbageTrigger: @@ -243,6 +245,8 @@ func (db *DB) triggerGarbageCollection() { // writeGCSizeDelay duration to avoid very frequent // database operations. func (db *DB) writeGCSizeWorker() { + defer close(db.writeGCSizeWorkerDone) + for { select { case <-db.writeGCSizeTrigger: diff --git a/swarm/storage/localstore/localstore.go b/swarm/storage/localstore/localstore.go index 7a9fb54f55..6e834f4c5b 100644 --- a/swarm/storage/localstore/localstore.go +++ b/swarm/storage/localstore/localstore.go @@ -107,6 +107,12 @@ type DB struct { // this channel is closed when close function is called // to terminate other goroutines 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. @@ -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 // to signal another event if it // is triggered during already running function - collectGarbageTrigger: make(chan struct{}, 1), - writeGCSizeTrigger: make(chan struct{}, 1), - close: make(chan struct{}), + collectGarbageTrigger: make(chan struct{}, 1), + writeGCSizeTrigger: make(chan struct{}, 1), + close: make(chan struct{}), + collectGarbageWorkerDone: make(chan struct{}), + writeGCSizeWorkerDone: make(chan struct{}), } if db.capacity <= 0 { 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) { close(db.close) 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 { log.Error("localstore: write gc size", "err", err) }