From a3c93ca4c46cb8b63afd8fc22619d4a0c43d49cf Mon Sep 17 00:00:00 2001 From: Anton Evangelatov Date: Mon, 5 Mar 2018 18:29:25 +0100 Subject: [PATCH] swarm/network, storage: Removal of chunk request from MemStore upon timeout or failure (#284) --- swarm/network/stream/common_test.go | 15 ++-- swarm/network/stream/delivery.go | 6 +- swarm/network/stream/intervals_test.go | 5 +- swarm/network/stream/syncer_test.go | 3 +- swarm/storage/dbapi.go | 17 ++-- swarm/storage/dpa.go | 4 +- swarm/storage/dpa_test.go | 8 +- swarm/storage/localstore.go | 41 ++++----- swarm/storage/netstore.go | 9 +- swarm/storage/netstore_test.go | 119 +++++++++++++++++++++++++ swarm/storage/types.go | 23 ++++- swarm/swarm.go | 9 +- 12 files changed, 205 insertions(+), 54 deletions(-) create mode 100644 swarm/storage/netstore_test.go diff --git a/swarm/network/stream/common_test.go b/swarm/network/stream/common_test.go index f7c767b352..0e6bbcaad5 100644 --- a/swarm/network/stream/common_test.go +++ b/swarm/network/stream/common_test.go @@ -39,11 +39,12 @@ import ( "github.com/ethereum/go-ethereum/swarm/network" "github.com/ethereum/go-ethereum/swarm/state" "github.com/ethereum/go-ethereum/swarm/storage" + colorable "github.com/mattn/go-colorable" ) var ( adapter = flag.String("adapter", "sim", "type of simulation: sim|socket|exec|docker") - loglevel = flag.Int("loglevel", 2, "verbosity of logs") + loglevel = flag.Int("loglevel", 4, "verbosity of logs") ) var ( @@ -63,8 +64,8 @@ func init() { // protocol when using the exec adapter adapters.RegisterServices(services) - log.Root().SetHandler(log.LvlFilterHandler(log.Lvl(*loglevel), log.StreamHandler(os.Stderr, log.TerminalFormat(false)))) - + log.PrintOrigins(true) + log.Root().SetHandler(log.LvlFilterHandler(log.Lvl(*loglevel), log.StreamHandler(colorable.NewColorableStderr(), log.TerminalFormat(true)))) } // NewStreamerService @@ -73,10 +74,10 @@ func NewStreamerService(ctx *adapters.ServiceContext) (node.Service, error) { addr := toAddr(id) kad := network.NewKademlia(addr.Over(), network.NewKadParams()) store := stores[id].(*storage.LocalStore) - db := storage.NewDBAPI(store) + netStore := storage.NewNetStore(store, nil) + db := storage.NewDBAPI(netStore) delivery := NewDelivery(kad, db) deliveries[id] = delivery - netStore := storage.NewNetStore(store, nil) r := NewRegistry(addr, delivery, netStore, state.NewMemStore(), defaultSkipCheck) RegisterSwarmSyncerServer(r, db) RegisterSwarmSyncerClient(r, db) @@ -105,7 +106,9 @@ func newStreamerTester(t *testing.T) (*p2ptest.ProtocolTester, *Registry, *stora return nil, nil, nil, removeDataDir, err } - db := storage.NewDBAPI(localStore) + netStore := storage.NewNetStore(localStore, nil) + db := storage.NewDBAPI(netStore) + delivery := NewDelivery(to, db) streamer := NewRegistry(addr, delivery, localStore, state.NewMemStore(), defaultSkipCheck) teardown := func() { diff --git a/swarm/network/stream/delivery.go b/swarm/network/stream/delivery.go index 7b44d090f8..c0aeb1630d 100644 --- a/swarm/network/stream/delivery.go +++ b/swarm/network/stream/delivery.go @@ -38,7 +38,6 @@ type Delivery struct { overlay network.Overlay receiveC chan *ChunkDeliveryMsg getPeer func(discover.NodeID) *Peer - quit chan struct{} } func NewDelivery(overlay network.Overlay, db *storage.DBAPI) *Delivery { @@ -139,6 +138,7 @@ func (d *Delivery) handleRetrieveRequestMsg(sp *Peer, req *RetrieveRequestMsg) e if created { if err := d.RequestFromPeers(chunk.Key[:], false, sp.ID()); err != nil { log.Warn("unable to forward chunk request", "peer", sp.ID(), "key", chunk.Key, "err", err) + chunk.SetErrored(true) return nil } } @@ -148,11 +148,11 @@ func (d *Delivery) handleRetrieveRequestMsg(sp *Peer, req *RetrieveRequestMsg) e select { case <-chunk.ReqC: - case <-d.quit: - return case <-t.C: + chunk.SetErrored(true) return } + chunk.SetErrored(false) if req.SkipCheck { err := sp.Deliver(chunk, s.priority) diff --git a/swarm/network/stream/intervals_test.go b/swarm/network/stream/intervals_test.go index fe90bb8c88..7614deef4d 100644 --- a/swarm/network/stream/intervals_test.go +++ b/swarm/network/stream/intervals_test.go @@ -47,10 +47,11 @@ func newIntervalsStreamerService(ctx *adapters.ServiceContext) (node.Service, er addr := toAddr(id) kad := network.NewKademlia(addr.Over(), network.NewKadParams()) store := stores[id].(*storage.LocalStore) - db := storage.NewDBAPI(store) + + netStore := storage.NewNetStore(store, nil) + db := storage.NewDBAPI(netStore) delivery := NewDelivery(kad, db) deliveries[id] = delivery - netStore := storage.NewNetStore(store, nil) r := NewRegistry(addr, delivery, netStore, state.NewMemStore(), defaultSkipCheck) r.RegisterClientFunc(externalStreamName, func(p *Peer, t []byte, live bool) (Client, error) { diff --git a/swarm/network/stream/syncer_test.go b/swarm/network/stream/syncer_test.go index 938c33d98c..9afe38adbb 100644 --- a/swarm/network/stream/syncer_test.go +++ b/swarm/network/stream/syncer_test.go @@ -108,7 +108,8 @@ func testSyncBetweenNodes(t *testing.T, nodes, conns, chunkCount int, skipCheck // create DBAPI-s for all nodes dbs := make([]*storage.DBAPI, nodes) for i := 0; i < nodes; i++ { - dbs[i] = storage.NewDBAPI(sim.Stores[i].(*storage.LocalStore)) + netStore := storage.NewNetStore(sim.Stores[i].(*storage.LocalStore), nil) + dbs[i] = storage.NewDBAPI(netStore) } // collect hashes in po 1 bin for each node diff --git a/swarm/storage/dbapi.go b/swarm/storage/dbapi.go index e0fac87464..565ef9f8fa 100644 --- a/swarm/storage/dbapi.go +++ b/swarm/storage/dbapi.go @@ -18,35 +18,34 @@ package storage // wrapper of db-s to provide mockable custom local chunk store access to syncer type DBAPI struct { - db *LDBStore - loc *LocalStore + ns *NetStore } -func NewDBAPI(loc *LocalStore) *DBAPI { - return &DBAPI{loc.DbStore.(*LDBStore), loc} +func NewDBAPI(ns *NetStore) *DBAPI { + return &DBAPI{ns: ns} } // to obtain the chunks from key or request db entry only func (self *DBAPI) Get(key Key) (*Chunk, error) { - return self.loc.Get(key) + return self.ns.localStore.Get(key) } // current storage counter of chunk db func (self *DBAPI) CurrentBucketStorageIndex(po uint8) uint64 { - return self.db.CurrentBucketStorageIndex(po) + return self.ns.localStore.DbStore.CurrentBucketStorageIndex(po) } // iteration storage counter and proximity order func (self *DBAPI) Iterator(from uint64, to uint64, po uint8, f func(Key, uint64) bool) error { - return self.db.SyncIterator(from, to, po, f) + return self.ns.localStore.DbStore.SyncIterator(from, to, po, f) } // to obtain the chunks from key or request db entry only func (self *DBAPI) GetOrCreateRequest(key Key) (*Chunk, bool) { - return self.loc.GetOrCreateRequest(key) + return self.ns.localStore.GetOrCreateRequest(key) } // to obtain the chunks from key or request db entry only func (self *DBAPI) Put(chunk *Chunk) { - self.loc.Put(chunk) + self.ns.localStore.Put(chunk) } diff --git a/swarm/storage/dpa.go b/swarm/storage/dpa.go index 2a28fbcd03..ecf215964b 100644 --- a/swarm/storage/dpa.go +++ b/swarm/storage/dpa.go @@ -77,8 +77,8 @@ func NewLocalDPA(datadir string, basekey []byte) (*DPA, error) { } return NewDPA(&LocalStore{ - NewMemStore(dbStore, singletonSwarmCacheCapacity), - dbStore, + memStore: NewMemStore(dbStore, singletonSwarmCacheCapacity), + DbStore: dbStore, }, NewChunkerParams()), nil } diff --git a/swarm/storage/dpa_test.go b/swarm/storage/dpa_test.go index 75f1448dc4..8b486fee21 100644 --- a/swarm/storage/dpa_test.go +++ b/swarm/storage/dpa_test.go @@ -36,8 +36,8 @@ func TestDPArandom(t *testing.T) { db.setCapacity(50000) memStore := NewMemStore(db, defaultCacheCapacity) localStore := &LocalStore{ - memStore, - db, + memStore: memStore, + DbStore: db, } chunker := NewTreeChunker(NewChunkerParams()) dpa := &DPA{ @@ -94,8 +94,8 @@ func TestDPA_capacity(t *testing.T) { db := tdb.LDBStore memStore := NewMemStore(db, 0) localStore := &LocalStore{ - memStore, - db, + memStore: memStore, + DbStore: db, } chunker := NewTreeChunker(NewChunkerParams()) dpa := &DPA{ diff --git a/swarm/storage/localstore.go b/swarm/storage/localstore.go index 02bcd2385e..1f38d518b6 100644 --- a/swarm/storage/localstore.go +++ b/swarm/storage/localstore.go @@ -19,6 +19,7 @@ package storage import ( "encoding/binary" "fmt" + "sync" "github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/metrics" @@ -47,8 +48,9 @@ func NewDefaultStoreParams() (self *StoreParams) { // LocalStore is a combination of inmemory db over a disk persisted db // implements a Get/Put with fallback (caching) logic using any 2 ChunkStores type LocalStore struct { - memStore ChunkStore - DbStore ChunkStore + memStore *MemStore + DbStore *LDBStore + mu sync.Mutex } // This constructor uses MemStore and DbStore as components @@ -63,20 +65,6 @@ func NewLocalStore(hash SwarmHasher, params *StoreParams, basekey []byte, mockSt }, nil } -func NewTestLocalStore(path string) (*LocalStore, error) { - basekey := make([]byte, 32) - hasher := MakeHashFunc("SHA3") - dbStore, err := NewLDBStore(path, hasher, singletonSwarmDbCapacity, func(k Key) (ret uint8) { return uint8(Proximity(basekey[:], k[:])) }) - if err != nil { - return nil, err - } - localStore := &LocalStore{ - memStore: NewMemStore(dbStore, singletonSwarmDbCapacity), - DbStore: dbStore, - } - return localStore, nil -} - func NewTestLocalStoreForAddr(path string, basekey []byte) (*LocalStore, error) { hasher := MakeHashFunc("SHA3") dbStore, err := NewLDBStore(path, hasher, singletonSwarmDbCapacity, func(k Key) (ret uint8) { return uint8(Proximity(basekey[:], k[:])) }) @@ -91,12 +79,15 @@ func NewTestLocalStoreForAddr(path string, basekey []byte) (*LocalStore, error) } func (self *LocalStore) CacheCounter() uint64 { - return uint64(self.memStore.(*MemStore).Counter()) + return uint64(self.memStore.Counter()) } // LocalStore is itself a chunk store // unsafe, in that the data is not integrity checked func (self *LocalStore) Put(chunk *Chunk) { + self.mu.Lock() + defer self.mu.Unlock() + chunk.Size = int64(binary.LittleEndian.Uint64(chunk.SData[0:8])) c := &Chunk{ Key: Key(append([]byte{}, chunk.Key...)), @@ -115,6 +106,13 @@ func (self *LocalStore) Put(chunk *Chunk) { // so additional timeout may be needed to wrap this call if // ChunkStores are remote and can have long latency func (self *LocalStore) Get(key Key) (chunk *Chunk, err error) { + self.mu.Lock() + defer self.mu.Unlock() + + return self.get(key) +} + +func (self *LocalStore) get(key Key) (chunk *Chunk, err error) { chunk, err = self.memStore.Get(key) if err == nil { if chunk.ReqC != nil { @@ -137,13 +135,16 @@ func (self *LocalStore) Get(key Key) (chunk *Chunk, err error) { // retrieve logic common for local and network chunk retrieval requests func (self *LocalStore) GetOrCreateRequest(key Key) (chunk *Chunk, created bool) { + self.mu.Lock() + defer self.mu.Unlock() + var err error - chunk, err = self.Get(key) - if err == nil { + chunk, err = self.get(key) + if err == nil && !chunk.GetErrored() { log.Trace(fmt.Sprintf("LocalStore.GetOrRetrieve: %v found locally", key)) return chunk, false } - if err == ErrFetching { + if err == ErrFetching && !chunk.GetErrored() { log.Trace(fmt.Sprintf("LocalStore.GetOrRetrieve: %v hit on an existing request %v", key, chunk.ReqC)) return chunk, false } diff --git a/swarm/storage/netstore.go b/swarm/storage/netstore.go index 51f476d687..4f1b6021a8 100644 --- a/swarm/storage/netstore.go +++ b/swarm/storage/netstore.go @@ -48,12 +48,16 @@ func (self *NetStore) Get(key Key) (chunk *Chunk, err error) { } else { var created bool chunk, created = self.localStore.GetOrCreateRequest(key) + if chunk.ReqC == nil { return chunk, nil } if created { - if err := self.retrieve(chunk); err != nil { + err := self.retrieve(chunk) + if err != nil { + // mark chunk request as failed so that we can retry it later + chunk.SetErrored(true) return nil, err } } @@ -64,9 +68,12 @@ func (self *NetStore) Get(key Key) (chunk *Chunk, err error) { select { case <-t.C: + // mark chunk request as failed so that we can retry + chunk.SetErrored(true) return nil, ErrChunkNotFound case <-chunk.ReqC: } + chunk.SetErrored(false) return chunk, nil } diff --git a/swarm/storage/netstore_test.go b/swarm/storage/netstore_test.go new file mode 100644 index 0000000000..ac43618e3e --- /dev/null +++ b/swarm/storage/netstore_test.go @@ -0,0 +1,119 @@ +// Copyright 2018 The go-ethereum Authors +// This file is part of the go-ethereum library. +// +// The go-ethereum library is free software: you can redistribute it and/or modify +// it under the terms of the GNU Lesser General Public License as published by +// the Free Software Foundation, either version 3 of the License, or +// (at your option) any later version. +// +// The go-ethereum library is distributed in the hope that it will be useful, +// but WITHOUT ANY WARRANTY; without even the implied warranty of +// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +// GNU Lesser General Public License for more details. +// +// You should have received a copy of the GNU Lesser General Public License +// along with the go-ethereum library. If not, see . + +package storage + +import ( + "encoding/hex" + "errors" + "io/ioutil" + "testing" + "time" + + "github.com/ethereum/go-ethereum/swarm/network" +) + +var ( + errUnknown = errors.New("unknown error") +) + +type mockRetrieve struct { + requests map[string]int +} + +func NewMockRetrieve() *mockRetrieve { + return &mockRetrieve{requests: make(map[string]int)} +} + +func newDummyChunk(key Key) *Chunk { + chunk := NewChunk(key, make(chan bool)) + chunk.SData = []byte{3, 4, 5} + chunk.Size = 3 + + return chunk +} + +func (m *mockRetrieve) retrieve(chunk *Chunk) error { + hkey := hex.EncodeToString(chunk.Key) + m.requests[hkey] += 1 + + // on second call return error + if m.requests[hkey] == 2 { + return errUnknown + } + + // on third call return data + if m.requests[hkey] == 3 { + *chunk = *newDummyChunk(chunk.Key) + go func() { + time.Sleep(50 * time.Millisecond) + close(chunk.ReqC) + }() + + return nil + } + + return nil +} + +func TestNetstoreFailedRequest(t *testing.T) { + searchTimeout = 100 * time.Millisecond + + // setup + addr := network.RandomAddr() // tested peers peer address + + // temp datadir + datadir, err := ioutil.TempDir("", "netstore") + if err != nil { + t.Fatal(err) + } + + localStore, err := NewTestLocalStoreForAddr(datadir, addr.Over()) + if err != nil { + t.Fatal(err) + } + + r := NewMockRetrieve() + netStore := NewNetStore(localStore, r.retrieve) + + // first call + key := Key{} + _, err = netStore.Get(key) + if err == nil || err != ErrChunkNotFound { + t.Fatalf("expected to get ErrChunkNotFound, but got: %s", err) + } + + // second call + _, err = netStore.Get(key) + if got := r.requests[hex.EncodeToString(key)]; got != 2 { + t.Fatalf("expected to have called retrieve two times, but got: %v", got) + } + if err != errUnknown { + t.Fatalf("expected to get an unknown error, but got: %s", err) + } + + // third call + chunk, err := netStore.Get(key) + if got := r.requests[hex.EncodeToString(key)]; got != 3 { + t.Fatalf("expected to have called retrieve three times, but got: %v", got) + } + if err != nil || chunk == nil { + t.Fatalf("expected to get a chunk but got: %v, %s", chunk, err) + } + if len(chunk.SData) != 3 { + t.Fatalf("expected to get a chunk with size 3, but got: %v", chunk.SData) + } +} diff --git a/swarm/storage/types.go b/swarm/storage/types.go index 2e6f6d7d47..7ea04322e1 100644 --- a/swarm/storage/types.go +++ b/swarm/storage/types.go @@ -24,6 +24,7 @@ import ( "fmt" "hash" "io" + "sync" "github.com/ethereum/go-ethereum/bmt" "github.com/ethereum/go-ethereum/common" @@ -173,9 +174,25 @@ type Chunk struct { SData []byte // nil if request, to be supplied by dpa Size int64 // size of the data covered by the subtree encoded in this chunk //Source Peer // peer - C chan bool // to signal data delivery by the dpa - ReqC chan bool // to signal the request done - dbStored chan bool // never remove a chunk from memStore before it is written to dbStore + C chan bool // to signal data delivery by the dpa + ReqC chan bool // to signal the request done + dbStored chan bool // never remove a chunk from memStore before it is written to dbStore + errored bool // flag which is set when the chunk request has errored or timeouted + erroredMu sync.Mutex +} + +func (c *Chunk) SetErrored(val bool) { + c.erroredMu.Lock() + defer c.erroredMu.Unlock() + + c.errored = val +} + +func (c *Chunk) GetErrored() bool { + c.erroredMu.Lock() + defer c.erroredMu.Unlock() + + return c.errored } func NewChunk(key Key, reqC chan bool) *Chunk { diff --git a/swarm/swarm.go b/swarm/swarm.go index 1fc554ecc0..d2c073c33c 100644 --- a/swarm/swarm.go +++ b/swarm/swarm.go @@ -146,7 +146,10 @@ func NewSwarm(ctx *node.ServiceContext, backend chequebook.Backend, config *api. HiveParams: config.HiveParams, } - db := storage.NewDBAPI(self.lstore) + // init netStore + ns := storage.NewNetStore(self.lstore, self.streamer.Retrieve) + + db := storage.NewDBAPI(ns) delivery := stream.NewDelivery(to, db) // TODO: decide on intervals store file location stateStore, err := state.NewDBStore(filepath.Join(config.Path, "state-store.db")) @@ -160,10 +163,10 @@ func NewSwarm(ctx *node.ServiceContext, backend chequebook.Backend, config *api. self.bzz = network.NewBzz(bzzconfig, to, stateStore, stream.Spec, self.streamer.Run) // set up DPA, the cloud storage local access layer - dpaChunkStore := storage.NewNetStore(self.lstore, self.streamer.Retrieve) + //dpaChunkStore := storage.NewNetStore(self.lstore, self.streamer.Retrieve) log.Debug(fmt.Sprintf("-> Local Access to Swarm")) // Swarm Hash Merklised Chunking for Arbitrary-length Document/File storage - self.dpa = storage.NewDPA(dpaChunkStore, self.config.ChunkerParams) + self.dpa = storage.NewDPA(ns, self.config.ChunkerParams) log.Debug(fmt.Sprintf("-> Content Store API")) // Pss = postal service over swarm (devp2p over bzz)