diff --git a/pstoreds/addr_book.go b/pstoreds/addr_book.go index 45c8202..ff2dfa4 100644 --- a/pstoreds/addr_book.go +++ b/pstoreds/addr_book.go @@ -118,20 +118,17 @@ func (r *addrsRecord) Clean() (chgd bool) { // dsAddrBook is an address book backed by a Datastore with a GC-like procedure // to purge expired entries. It uses an in-memory address stream manager. type dsAddrBook struct { - ctx context.Context - + ctx context.Context opts Options cache cache ds ds.Batching + gc *dsAddrBookGc subsManager *pstoremem.AddrSubManager flushJobCh chan *addrsRecord cancelFn func() - closeDone sync.WaitGroup - - gcCurrWindowEnd int64 - gcLookaheadRunning int32 + done sync.WaitGroup } var _ pstore.AddrBook = (*dsAddrBook)(nil) @@ -157,16 +154,18 @@ func NewAddrBook(ctx context.Context, store ds.Batching, opts Options) (ab *dsAd flushJobCh: make(chan *addrsRecord, 32), } + ab.gc = newAddressBookGc(ctx, ab) + // kick off background processes. go ab.flusher() - go ab.gc() + go ab.gc.background() return ab, nil } func (ab *dsAddrBook) Close() { ab.cancelFn() - ab.closeDone.Wait() + ab.done.Wait() } func (ab *dsAddrBook) asyncFlush(pr *addrsRecord) { @@ -219,7 +218,9 @@ func (ab *dsAddrBook) loadRecord(id peer.ID, cache bool, update bool) (pr *addrs // flusher is a goroutine that takes care of persisting asynchronous flushes to the datastore. func (ab *dsAddrBook) flusher() { - ab.closeDone.Add(1) + ab.done.Add(1) + defer ab.done.Done() + for { select { case fj := <-ab.flushJobCh: @@ -236,37 +237,6 @@ func (ab *dsAddrBook) flusher() { } case <-ab.ctx.Done(): - ab.closeDone.Done() - return - } - } -} - -// gc is a goroutine that prunes expired addresses from the datastore at regular intervals. -func (ab *dsAddrBook) gc() { - select { - case <-time.After(ab.opts.GCInitialDelay): - case <-ab.ctx.Done(): - // yield if we have been cancelled/closed before the delay elapses. - return - } - - ab.closeDone.Add(1) - purgeTimer := time.NewTicker(ab.opts.GCPurgeInterval) - lookaheadTimer := time.NewTicker(ab.opts.GCLookaheadInterval) - - for { - select { - case <-purgeTimer.C: - ab.purgeCycle() - - case <-lookaheadTimer.C: - ab.populateLookahead() - - case <-ab.ctx.Done(): - purgeTimer.Stop() - lookaheadTimer.Stop() - ab.closeDone.Done() return } } diff --git a/pstoreds/addr_book_gc.go b/pstoreds/addr_book_gc.go index f175831..80a7a61 100644 --- a/pstoreds/addr_book_gc.go +++ b/pstoreds/addr_book_gc.go @@ -1,14 +1,13 @@ package pstoreds import ( + "context" "fmt" "strconv" - "sync/atomic" "time" ds "github.com/ipfs/go-datastore" query "github.com/ipfs/go-datastore/query" - errors "github.com/pkg/errors" peer "github.com/libp2p/go-libp2p-peer" pb "github.com/libp2p/go-libp2p-peerstore/pb" @@ -20,14 +19,17 @@ var ( // GC lookahead entries are stored in keys with pattern: // /peers/gc/addrs// => nil gcLookaheadBase = ds.NewKey("/peers/gc/addrs") + // in GC routines, how many operations do we place in a batch before it's committed. gcOpsPerBatch = 20 + // queries purgeQuery = query.Query{ Prefix: gcLookaheadBase.String(), Orders: []query.Order{query.OrderByKey{}}, KeysOnly: true, } + populateLookaheadQuery = query.Query{ Prefix: addrBookBase.String(), Orders: []query.Order{query.OrderByKey{}}, @@ -35,82 +37,73 @@ var ( } ) -// cyclicBatch buffers datastore write operations and automatically flushes them after gcOpsPerBatch (20) have been -// queued. An explicit `Commit()` closes this cyclic batch, erroring all further operations. -// -// It is similar to go-datastore autobatch, but it's driven by an actual Batch facility offered by the -// datastore. -type cyclicBatch struct { - ds.Batch - ds ds.Batching - pending int +// dsAddrBookGc encapsulates the GC behaviour to maintain a datastore-backed address book. +type dsAddrBookGc struct { + ctx context.Context + ab *dsAddrBook + running chan struct{} + currWindowEnd int64 } -func newCyclicBatch(ds ds.Batching) (ds.Batch, error) { - batch, err := ds.Batch() - if err != nil { - return nil, err +func newAddressBookGc(ctx context.Context, ab *dsAddrBook) *dsAddrBookGc { + return &dsAddrBookGc{ + ctx: ctx, + ab: ab, + running: make(chan struct{}, 1), } - return &cyclicBatch{Batch: batch, ds: ds}, nil } -func (cb *cyclicBatch) cycle() (err error) { - if cb.Batch == nil { - return errors.New("cyclic batch is closed") - } - if cb.pending < gcOpsPerBatch { - // we haven't reached the threshold yet. - return nil - } - // commit and renew the batch. - if err = cb.Batch.Commit(); err != nil { - return errors.Wrap(err, "failed while committing cyclic batch") - } - if cb.Batch, err = cb.ds.Batch(); err != nil { - return errors.Wrap(err, "failed while renewing cyclic batch") - } - return nil -} +// gc prunes expired addresses from the datastore at regular intervals. It should be spawned as a goroutine. +func (gc *dsAddrBookGc) background() { + gc.ab.done.Add(1) + defer gc.ab.done.Done() -func (cb *cyclicBatch) Put(key ds.Key, val []byte) error { - if err := cb.cycle(); err != nil { - return err + select { + case <-time.After(gc.ab.opts.GCInitialDelay): + case <-gc.ab.ctx.Done(): + // yield if we have been cancelled/closed before the delay elapses. + return } - cb.pending++ - return cb.Batch.Put(key, val) -} -func (cb *cyclicBatch) Delete(key ds.Key) error { - if err := cb.cycle(); err != nil { - return err - } - cb.pending++ - return cb.Batch.Delete(key) -} + purgeTimer := time.NewTicker(gc.ab.opts.GCPurgeInterval) + defer purgeTimer.Stop() -func (cb *cyclicBatch) Commit() error { - if cb.Batch == nil { - return errors.New("cyclic batch is closed") + var lookaheadCh <-chan time.Time + if gc.ab.opts.GCLookaheadInterval > 0 { + lookaheadTimer := time.NewTicker(gc.ab.opts.GCLookaheadInterval) + lookaheadCh = lookaheadTimer.C + defer lookaheadTimer.Stop() } - if err := cb.Batch.Commit(); err != nil { - return err + + for { + select { + case <-purgeTimer.C: + gc.purgeCycle() + + case <-lookaheadCh: + // will never trigger if lookahead is disabled (nil Duration). + gc.populateLookahead() + + case <-gc.ctx.Done(): + return + } } - cb.pending = 0 - cb.Batch = nil - return nil } // purgeCycle runs a single GC purge cycle. It operates within the lookahead window if lookahead is enabled; else it // visits all entries in the datastore, deleting the addresses that have expired. -func (ab *dsAddrBook) purgeCycle() { - if atomic.LoadInt32(&ab.gcLookaheadRunning) > 0 { +func (gc *dsAddrBookGc) purgeCycle() { + select { + case gc.running <- struct{}{}: + defer func() { <-gc.running }() + default: // yield if lookahead is running. return } var id peer.ID record := &addrsRecord{AddrBookRecord: &pb.AddrBookRecord{}} // empty record to reuse and avoid allocs. - batch, err := newCyclicBatch(ab.ds) + batch, err := newCyclicBatch(gc.ab.ds) if err != nil { log.Warningf("failed while creating batch to purge GC entries: %v", err) } @@ -135,7 +128,7 @@ func (ab *dsAddrBook) purgeCycle() { } // re-add the record if it needs to be visited again in this window. - if len(ar.Addrs) != 0 && ar.Addrs[0].Expiry <= ab.gcCurrWindowEnd { + if len(ar.Addrs) != 0 && ar.Addrs[0].Expiry <= gc.currWindowEnd { gcKey := gcLookaheadBase.ChildString(fmt.Sprintf("%d/%s", ar.Addrs[0].Expiry, key.Name())) if err := batch.Put(gcKey, []byte{}); err != nil { log.Warningf("failed to add new GC key: %v, err: %v", gcKey, err) @@ -143,7 +136,7 @@ func (ab *dsAddrBook) purgeCycle() { } } - results, err := ab.ds.Query(purgeQuery) + results, err := gc.ab.ds.Query(purgeQuery) if err != nil { log.Warningf("failed while fetching entries to purge: %v", err) return @@ -181,7 +174,7 @@ func (ab *dsAddrBook) purgeCycle() { } // if the record is in cache, we clean it and flush it if necessary. - if e, ok := ab.cache.Peek(id); ok { + if e, ok := gc.ab.cache.Peek(id); ok { cached := e.(*addrsRecord) cached.Lock() if cached.Clean() { @@ -198,7 +191,7 @@ func (ab *dsAddrBook) purgeCycle() { // otherwise, fetch it from the store, clean it and flush it. entryKey := addrBookBase.ChildString(gcKey.Name()) - val, err := ab.ds.Get(entryKey) + val, err := gc.ab.ds.Get(entryKey) if err != nil { // captures all errors, including ErrNotFound. dropInError(gcKey, err, "fetching entry") @@ -224,27 +217,31 @@ func (ab *dsAddrBook) purgeCycle() { } // populateLookahead populates the lookahead window by scanning the entire store and picking entries whose earliest -// expiration falls within the new window. +// expiration falls within the window period. // // Those entries are stored in the lookahead region in the store, indexed by the timestamp when they need to be // visited, to facilitate temporal range scans. -func (ab *dsAddrBook) populateLookahead() { - if !atomic.CompareAndSwapInt32(&ab.gcLookaheadRunning, 0, 1) { +func (gc *dsAddrBookGc) populateLookahead() { + select { + case gc.running <- struct{}{}: + defer func() { <-gc.running }() + default: + // yield if something's running. return } - until := time.Now().Add(ab.opts.GCLookaheadInterval).Unix() + until := time.Now().Add(gc.ab.opts.GCLookaheadInterval).Unix() var id peer.ID record := &addrsRecord{AddrBookRecord: &pb.AddrBookRecord{}} - results, err := ab.ds.Query(populateLookaheadQuery) + results, err := gc.ab.ds.Query(populateLookaheadQuery) if err != nil { log.Warningf("failed while querying to populate lookahead GC window: %v", err) return } defer results.Close() - batch, err := newCyclicBatch(ab.ds) + batch, err := newCyclicBatch(gc.ab.ds) if err != nil { log.Warningf("failed while creating batch to populate lookahead GC window: %v", err) return @@ -262,7 +259,7 @@ func (ab *dsAddrBook) populateLookahead() { } // if the record is in cache, use the cached version. - if e, ok := ab.cache.Peek(id); ok { + if e, ok := gc.ab.cache.Peek(id); ok { cached := e.(*addrsRecord) cached.RLock() if len(cached.Addrs) == 0 || cached.Addrs[0].Expiry > until { @@ -279,7 +276,7 @@ func (ab *dsAddrBook) populateLookahead() { record.Reset() - val, err := ab.ds.Get(ds.RawKey(result.Key)) + val, err := gc.ab.ds.Get(ds.RawKey(result.Key)) if err != nil { log.Warningf("failed which getting record from store for peer: %v, err: %v", id.Pretty(), err) continue @@ -300,6 +297,5 @@ func (ab *dsAddrBook) populateLookahead() { log.Warningf("failed to commit GC lookahead batch: %v", err) } - ab.gcCurrWindowEnd = until - atomic.StoreInt32(&ab.gcLookaheadRunning, 0) + gc.currWindowEnd = until } diff --git a/pstoreds/addr_book_gc_test.go b/pstoreds/addr_book_gc_test.go index 509817f..3d43363 100644 --- a/pstoreds/addr_book_gc_test.go +++ b/pstoreds/addr_book_gc_test.go @@ -44,6 +44,7 @@ func TestGCLookahead(t *testing.T) { factory := addressBookFactory(t, badgerStore, opts) ab, closeFn := factory() + gc := ab.(*dsAddrBook).gc defer closeFn() tp := &testProbe{t, ab} @@ -55,7 +56,8 @@ func TestGCLookahead(t *testing.T) { ab.AddAddrs(ids[0], addrs[:10], time.Hour) ab.AddAddrs(ids[1], addrs[10:20], time.Hour) ab.AddAddrs(ids[2], addrs[20:30], time.Hour) - ab.(*dsAddrBook).populateLookahead() + + gc.populateLookahead() if i := tp.countLookaheadEntries(); i != 0 { t.Errorf("expected no GC lookahead entries, got: %v", i) } @@ -66,14 +68,14 @@ func TestGCLookahead(t *testing.T) { // Purge the cache, to exercise a different path in the lookahead cycle. tp.clearCache() - ab.(*dsAddrBook).populateLookahead() + gc.populateLookahead() if i := tp.countLookaheadEntries(); i != 1 { t.Errorf("expected 1 GC lookahead entry, got: %v", i) } // change addresses of another to have TTL 5 second, placing them in the lookahead window. ab.UpdateAddrs(ids[2], time.Hour, 5*time.Second) - ab.(*dsAddrBook).populateLookahead() + gc.populateLookahead() if i := tp.countLookaheadEntries(); i != 2 { t.Errorf("expected 2 GC lookahead entries, got: %v", i) } @@ -89,6 +91,7 @@ func TestGCPurging(t *testing.T) { factory := addressBookFactory(t, badgerStore, opts) ab, closeFn := factory() + gc := ab.(*dsAddrBook).gc defer closeFn() tp := &testProbe{t, ab} @@ -110,13 +113,13 @@ func TestGCPurging(t *testing.T) { // this is inside the window, but it will survive the purges we do in the test. ab.AddAddrs(ids[3], addrs[70:80], 15*time.Second) - ab.(*dsAddrBook).populateLookahead() + gc.populateLookahead() if i := tp.countLookaheadEntries(); i != 4 { t.Errorf("expected 4 GC lookahead entries, got: %v", i) } <-time.After(2 * time.Second) - ab.(*dsAddrBook).purgeCycle() + gc.purgeCycle() if i := tp.countLookaheadEntries(); i != 3 { t.Errorf("expected 3 GC lookahead entries, got: %v", i) } @@ -125,13 +128,13 @@ func TestGCPurging(t *testing.T) { tp.clearCache() <-time.After(5 * time.Second) - ab.(*dsAddrBook).purgeCycle() + gc.purgeCycle() if i := tp.countLookaheadEntries(); i != 3 { t.Errorf("expected 3 GC lookahead entries, got: %v", i) } <-time.After(5 * time.Second) - ab.(*dsAddrBook).purgeCycle() + gc.purgeCycle() if i := tp.countLookaheadEntries(); i != 1 { t.Errorf("expected 1 GC lookahead entries, got: %v", i) } @@ -200,6 +203,6 @@ func BenchmarkLookaheadCycle(b *testing.B) { b.ResetTimer() for i := 0; i < b.N; i++ { - ab.(*dsAddrBook).populateLookahead() + ab.(*dsAddrBook).gc.populateLookahead() } } diff --git a/pstoreds/cyclic_batch.go b/pstoreds/cyclic_batch.go new file mode 100644 index 0000000..64dc103 --- /dev/null +++ b/pstoreds/cyclic_batch.go @@ -0,0 +1,72 @@ +package pstoreds + +import ( + "github.com/pkg/errors" + + ds "github.com/ipfs/go-datastore" +) + +// cyclicBatch buffers ds write operations and automatically flushes them after gcOpsPerBatch (20) have been +// queued. An explicit `Commit()` closes this cyclic batch, erroring all further operations. +// +// It is similar to go-ds autobatch, but it's driven by an actual Batch facility offered by the +// ds. +type cyclicBatch struct { + ds.Batch + ds ds.Batching + pending int +} + +func newCyclicBatch(ds ds.Batching) (ds.Batch, error) { + batch, err := ds.Batch() + if err != nil { + return nil, err + } + return &cyclicBatch{Batch: batch, ds: ds}, nil +} + +func (cb *cyclicBatch) cycle() (err error) { + if cb.Batch == nil { + return errors.New("cyclic batch is closed") + } + if cb.pending < gcOpsPerBatch { + // we haven't reached the threshold yet. + return nil + } + // commit and renew the batch. + if err = cb.Batch.Commit(); err != nil { + return errors.Wrap(err, "failed while committing cyclic batch") + } + if cb.Batch, err = cb.ds.Batch(); err != nil { + return errors.Wrap(err, "failed while renewing cyclic batch") + } + return nil +} + +func (cb *cyclicBatch) Put(key ds.Key, val []byte) error { + if err := cb.cycle(); err != nil { + return err + } + cb.pending++ + return cb.Batch.Put(key, val) +} + +func (cb *cyclicBatch) Delete(key ds.Key) error { + if err := cb.cycle(); err != nil { + return err + } + cb.pending++ + return cb.Batch.Delete(key) +} + +func (cb *cyclicBatch) Commit() error { + if cb.Batch == nil { + return errors.New("cyclic batch is closed") + } + if err := cb.Batch.Commit(); err != nil { + return err + } + cb.pending = 0 + cb.Batch = nil + return nil +}