From a15cb44f986c8ac383eec49ffc465c142e46e336 Mon Sep 17 00:00:00 2001 From: lash Date: Tue, 9 Oct 2018 18:40:51 +0200 Subject: [PATCH] swarm/storage: Use main batch for GC, correct locking (locking now works but implementation is messy) --- swarm/storage/ldbstore.go | 94 ++++++++++++++++++++-------------- swarm/storage/ldbstore_test.go | 48 +++++++++++++---- 2 files changed, 94 insertions(+), 48 deletions(-) diff --git a/swarm/storage/ldbstore.go b/swarm/storage/ldbstore.go index 8c8b7a3ddf..7ab3f9cb5b 100644 --- a/swarm/storage/ldbstore.go +++ b/swarm/storage/ldbstore.go @@ -97,7 +97,9 @@ type garbage struct { target int // number of chunks to delete in running round batch *dbBatch // the delete batch - wg sync.WaitGroup // set to wait when a gc round is active + running bool + + //wg sync.WaitGroup // set to wait when a gc round is active } type LDBStore struct { @@ -182,7 +184,7 @@ func NewLDBStore(params *LDBStoreParams) (s *LDBStore, err error) { maxBatch: defaultMaxGCBatch, maxRound: defaultMaxGCRound, ratio: defaultGCRatio, - batch: newBatch(), + //batch: newBatch(), } return s, nil @@ -202,14 +204,14 @@ func (s *LDBStore) startGC(c int) { } // commit deletions to db -func (s *LDBStore) runGC() error { - err := s.db.Write(s.gc.batch.Batch) - if err != nil { - return err - } - s.gc.batch.Reset() - return nil -} +//func (s *LDBStore) runGC() error { +// err := s.db.Write(s.gc.batch.Batch) +// if err != nil { +// return err +// } +// s.gc.batch.Reset() +// return nil +//} // NewMockDbStore creates a new instance of DbStore with // mockStore set to a provided value. If mockStore argument is nil, @@ -316,10 +318,18 @@ func decodeData(addr Address, data []byte) (*chunk, error) { func (s *LDBStore) collectGarbage() { - metrics.GetOrRegisterCounter("ldbstore.collectgarbage", nil).Inc(1) + s.lock.Lock() + if s.gc.running { + s.lock.Unlock() + return + } + s.gc.running = true + defer func() { + s.gc.running = false + }() + s.lock.Unlock() - s.gc.wg.Add(1) - defer s.gc.wg.Done() + metrics.GetOrRegisterCounter("ldbstore.collectgarbage", nil).Inc(1) it := s.db.NewIterator() defer it.Release() @@ -331,6 +341,7 @@ func (s *LDBStore) collectGarbage() { ok := it.Seek([]byte{keyGCIdx}) for s.gc.count < s.gc.target { var singleIterationCount int + s.lock.Lock() for ; ok && (singleIterationCount < s.gc.maxBatch); ok = it.Next() { itkey := it.Key() @@ -346,17 +357,19 @@ func (s *LDBStore) collectGarbage() { log.Trace("parse gc", "index", index, "po", po, "hash", hash) - s.delete(s.gc.batch.Batch, index, keyIdx, po) + s.delete(s.batch.Batch, index, keyIdx, po) singleIterationCount++ s.gc.count++ if s.gc.count > s.gc.maxRound { break } } - err := s.runGC() - if err != nil { - log.Error("gc fail: %v", err) - } + s.lock.Unlock() + s.batchesC <- struct{}{} + // err := s.runGC() + // if err != nil { + // log.Error("gc fail: %v", err) + // } log.Trace("garbage collect batch done", "batch", singleIterationCount, "total", s.gc.count) } log.Debug("garbage collect done", "c", s.gc.count) @@ -650,13 +663,12 @@ func (s *LDBStore) Put(ctx context.Context, chunk Chunk) error { metrics.GetOrRegisterCounter("ldbstore.put", nil).Inc(1) log.Trace("ldbstore.put", "key", chunk.Address()) + s.lock.Lock() ikey := getIndexKey(chunk.Address()) var index dpaDBIndex po := s.po(chunk.Address()) - s.lock.Lock() - if s.closed { s.lock.Unlock() return ErrDBClosed @@ -693,6 +705,7 @@ func (s *LDBStore) Put(ctx context.Context, chunk Chunk) error { case <-ctx.Done(): return ctx.Err() } + } // force putting into db, does not check access index @@ -732,10 +745,10 @@ func (s *LDBStore) writeBatches() { func (s *LDBStore) writeCurrentBatch() error { s.lock.Lock() - defer s.lock.Unlock() b := s.batch l := b.Len() if l == 0 { + s.lock.Unlock() return nil } e := s.entryCnt @@ -744,24 +757,27 @@ func (s *LDBStore) writeCurrentBatch() error { s.batch = newBatch() b.err = s.writeBatch(b, e, d, a) close(b.c) - for e > s.capacity { - log.Debug("for >", "e", e, "s.capacity", s.capacity) - // Collect garbage in a separate goroutine - // to be able to interrupt this loop by s.quit. - done := make(chan struct{}) - go func() { - s.collectGarbage() - log.Trace("collectGarbage closing done") - close(done) - }() - - select { - case <-s.quit: - return errors.New("CollectGarbage terminated due to quit") - case <-done: - } - e = s.entryCnt + s.lock.Unlock() + if e > s.capacity { + go s.collectGarbage() } + // log.Debug("for >", "e", e, "s.capacity", s.capacity) + // // Collect garbage in a separate goroutine + // // to be able to interrupt this loop by s.quit. + // done := make(chan struct{}) + // go func() { + // s.collectGarbage() + // log.Trace("collectGarbage closing done") + // close(done) + // }() + // + // select { + // case <-s.quit: + // return errors.New("CollectGarbage terminated due to quit") + // case <-done: + // } + // e = s.entryCnt + // } return nil } @@ -793,6 +809,8 @@ func newMockEncodeDataFunc(mockStore *mock.NodeStore) func(chunk Chunk) []byte { // try to find index; if found, update access cnt and return true func (s *LDBStore) tryAccessIdx(ikey []byte, po uint8, index *dpaDBIndex) bool { + s.lock.Lock() + defer s.lock.Unlock() idata, err := s.db.Get(ikey) if err != nil { return false diff --git a/swarm/storage/ldbstore_test.go b/swarm/storage/ldbstore_test.go index 5655412b1a..f6f71caa42 100644 --- a/swarm/storage/ldbstore_test.go +++ b/swarm/storage/ldbstore_test.go @@ -19,6 +19,7 @@ package storage import ( "bytes" "context" + "errors" "fmt" "io/ioutil" "os" @@ -279,7 +280,7 @@ func TestLDBStoreWithoutCollectGarbage(t *testing.T) { log.Info("ldbstore", "entrycnt", ldb.entryCnt, "accesscnt", ldb.accessCnt) for _, ch := range chunks { - ret, err := ldb.Get(context.TODO(), ch.Address()) + ret, err := ldb.get(ch.Address()) if err != nil { t.Fatal(err) } @@ -338,11 +339,13 @@ func testLDBStoreCollectGarbage(t *testing.T) { log.Info("ldbstore", "entrycnt", ldb.entryCnt, "accesscnt", ldb.accessCnt) // wait for garbage collection to kick in on the responsible actor - ldb.gc.wg.Wait() + ctx, cancel := context.WithTimeout(context.Background(), time.Second*10) + defer cancel() + waitGc(ctx, ldb) var missing int for _, ch := range chunks { - ret, err := ldb.Get(context.Background(), ch.Address()) + ret, err := ldb.get(ch.Address()) if err == ErrChunkNotFound || err == ldberrors.ErrNotFound { missing++ continue @@ -439,7 +442,13 @@ func testLDBStoreRemoveThenCollectGarbage(t *testing.T) { // delete all chunks for i := 0; i < n; i++ { - ldb.Delete(chunks[i].Address()) + ikey := getIndexKey(chunks[i].Address()) + + var indx dpaDBIndex + proximity := ldb.po(chunks[i].Address()) + ldb.tryAccessIdx(ikey, proximity, &indx) + + ldb.deleteNow(&indx, ikey, proximity) } log.Info("ldbstore", "entrycnt", ldb.entryCnt, "accesscnt", ldb.accessCnt) @@ -466,11 +475,13 @@ func testLDBStoreRemoveThenCollectGarbage(t *testing.T) { } // wait for garbage collection - time.Sleep(1 * time.Second) + ctx, cancel := context.WithTimeout(context.Background(), time.Second*10) + defer cancel() + waitGc(ctx, ldb) // expect first surplus chunks to be missing, because they have the smallest access value for i := 0; i < surplus; i++ { - _, err := ldb.Get(context.TODO(), chunks[i].Address()) + _, err := ldb.get(chunks[i].Address()) if err == nil { t.Fatal("expected surplus chunk to be missing, but got no error") } @@ -478,7 +489,7 @@ func testLDBStoreRemoveThenCollectGarbage(t *testing.T) { // expect last chunks to be present, as they have the largest access value for i := surplus; i < surplus+capacity; i++ { - ret, err := ldb.Get(context.TODO(), chunks[i].Address()) + ret, err := ldb.get(chunks[i].Address()) if err != nil { t.Fatalf("chunk %v: expected no error, but got %s", i, err) } @@ -506,7 +517,7 @@ func TestLDBStoreCollectGarbageAccessUnlikeIndex(t *testing.T) { // set first added capacity/2 chunks to highest accesscount for i := 0; i < capacity/2; i++ { - _, err := ldb.Get(context.TODO(), chunks[i].Address()) + _, err := ldb.get(chunks[i].Address()) if err != nil { t.Fatalf("fail add chunk #%d - %s: %v", i, chunks[i].Address(), err) } @@ -517,11 +528,13 @@ func TestLDBStoreCollectGarbageAccessUnlikeIndex(t *testing.T) { } // wait for garbage collection to kick in on the responsible actor - ldb.gc.wg.Wait() + ctx, cancel := context.WithTimeout(context.Background(), time.Second*10) + defer cancel() + waitGc(ctx, ldb) var missing int for i, ch := range chunks[2 : capacity/2] { - ret, err := ldb.Get(context.Background(), ch.Address()) + ret, err := ldb.get(ch.Address()) if err == ErrChunkNotFound || err == ldberrors.ErrNotFound { t.Fatalf("fail find chunk #%d - %s: %v", i, ch.Address(), err) } @@ -534,3 +547,18 @@ func TestLDBStoreCollectGarbageAccessUnlikeIndex(t *testing.T) { log.Info("ldbstore", "total", n, "missing", missing, "entrycnt", ldb.entryCnt, "accesscnt", ldb.accessCnt) } + +func waitGc(ctx context.Context, ldb *LDBStore) error { + ticker := time.Tick(time.Millisecond * 100) + for { + select { + + case <-ctx.Done(): + return errors.New("timeout") + case <-ticker: + if !ldb.gc.running { + return nil + } + } + } +}