From ad1ff36aaf7323497ae09d2bf08383e7d7951b77 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jano=C5=A1=20Gulja=C5=A1?= Date: Mon, 5 Mar 2018 18:42:47 +0100 Subject: [PATCH] swarm/network/stream: fix sync tests (#297) --- swarm/network/stream/common_test.go | 21 ++++- swarm/network/stream/intervals_test.go | 2 +- swarm/network/stream/stream.go | 120 ++++++++++++------------- swarm/network/stream/syncer_test.go | 1 + swarm/swarm.go | 4 +- 5 files changed, 79 insertions(+), 69 deletions(-) diff --git a/swarm/network/stream/common_test.go b/swarm/network/stream/common_test.go index 0e6bbcaad5..c14a078087 100644 --- a/swarm/network/stream/common_test.go +++ b/swarm/network/stream/common_test.go @@ -32,6 +32,7 @@ import ( "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/node" + "github.com/ethereum/go-ethereum/p2p" "github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/p2p/simulations/adapters" p2ptest "github.com/ethereum/go-ethereum/p2p/testing" @@ -78,13 +79,14 @@ func NewStreamerService(ctx *adapters.ServiceContext) (node.Service, error) { db := storage.NewDBAPI(netStore) delivery := NewDelivery(kad, db) deliveries[id] = delivery - r := NewRegistry(addr, delivery, netStore, state.NewMemStore(), defaultSkipCheck) + r := NewRegistry(addr, delivery, db, state.NewMemStore(), defaultSkipCheck, false) RegisterSwarmSyncerServer(r, db) RegisterSwarmSyncerClient(r, db) go func() { waitPeerErrC <- waitForPeers(r, 1*time.Second, peerCount(id)) }() - return &TestRegistry{Registry: r}, nil + dpa := storage.NewDPA(storage.NewNetStore(store, nil), storage.NewChunkerParams()) + return &TestRegistry{Registry: r, dpa: dpa}, nil } func newStreamerTester(t *testing.T) (*p2ptest.ProtocolTester, *Registry, *storage.LocalStore, func(), error) { @@ -110,7 +112,7 @@ func newStreamerTester(t *testing.T) (*p2ptest.ProtocolTester, *Registry, *stora db := storage.NewDBAPI(netStore) delivery := NewDelivery(to, db) - streamer := NewRegistry(addr, delivery, localStore, state.NewMemStore(), defaultSkipCheck) + streamer := NewRegistry(addr, delivery, db, state.NewMemStore(), defaultSkipCheck, false) teardown := func() { streamer.Close() removeDataDir() @@ -169,6 +171,7 @@ func (rrs *roundRobinStore) Close() { type TestRegistry struct { *Registry + dpa *storage.DPA } func (r *TestRegistry) APIs() []rpc.API { @@ -199,7 +202,17 @@ func readAll(dpa *storage.DPA, hash []byte) (int64, error) { } func (r *TestRegistry) ReadAll(hash common.Hash) (int64, error) { - return readAll(r.api.dpa, hash[:]) + return readAll(r.dpa, hash[:]) +} + +func (r *TestRegistry) Start(server *p2p.Server) error { + r.dpa.Start() + return r.Registry.Start(server) +} + +func (r *TestRegistry) Stop() error { + r.dpa.Stop() + return r.Registry.Stop() } type TestExternalRegistry struct { diff --git a/swarm/network/stream/intervals_test.go b/swarm/network/stream/intervals_test.go index 7614deef4d..af6496fb73 100644 --- a/swarm/network/stream/intervals_test.go +++ b/swarm/network/stream/intervals_test.go @@ -52,7 +52,7 @@ func newIntervalsStreamerService(ctx *adapters.ServiceContext) (node.Service, er db := storage.NewDBAPI(netStore) delivery := NewDelivery(kad, db) deliveries[id] = delivery - r := NewRegistry(addr, delivery, netStore, state.NewMemStore(), defaultSkipCheck) + r := NewRegistry(addr, delivery, db, state.NewMemStore(), defaultSkipCheck, false) r.RegisterClientFunc(externalStreamName, func(p *Peer, t []byte, live bool) (Client, error) { return newTestExternalClient(t, db), nil diff --git a/swarm/network/stream/stream.go b/swarm/network/stream/stream.go index 0b8c23b362..a145b6f180 100644 --- a/swarm/network/stream/stream.go +++ b/swarm/network/stream/stream.go @@ -56,23 +56,23 @@ type Registry struct { clientFuncs map[string]func(*Peer, []byte, bool) (Client, error) peers map[discover.NodeID]*Peer delivery *Delivery - store storage.ChunkStore intervalsStore state.Store + doSync bool } // NewRegistry is Streamer constructor -func NewRegistry(addr *network.BzzAddr, delivery *Delivery, store storage.ChunkStore, intervalsStore state.Store, skipCheck bool) *Registry { +func NewRegistry(addr *network.BzzAddr, delivery *Delivery, db *storage.DBAPI, intervalsStore state.Store, skipCheck, doSync bool) *Registry { streamer := &Registry{ addr: addr, skipCheck: skipCheck, - store: store, serverFuncs: make(map[string]func(*Peer, []byte, bool) (Server, error)), clientFuncs: make(map[string]func(*Peer, []byte, bool) (Client, error)), peers: make(map[discover.NodeID]*Peer), delivery: delivery, intervalsStore: intervalsStore, + doSync: doSync, } - streamer.api = NewAPI(streamer, streamer.store) + streamer.api = NewAPI(streamer) delivery.getPeer = streamer.getPeer streamer.RegisterServerFunc(swarmChunkServerStreamName, func(_ *Peer, _ []byte, _ bool) (Server, error) { return NewSwarmChunkServer(delivery.db), nil @@ -80,6 +80,8 @@ func NewRegistry(addr *network.BzzAddr, delivery *Delivery, store storage.ChunkS streamer.RegisterClientFunc(swarmChunkServerStreamName, func(p *Peer, _ []byte, _ bool) (Client, error) { return NewSwarmSyncerClient(p, delivery.db, nil) }) + RegisterSwarmSyncerServer(streamer, db) + RegisterSwarmSyncerClient(streamer, db) return streamer } @@ -214,7 +216,6 @@ func (r *Registry) PeerInfo(id discover.NodeID) interface{} { } func (r *Registry) Close() error { - r.store.Close() return r.intervalsStore.Close() } @@ -252,63 +253,65 @@ func (r *Registry) Run(p *network.BzzPeer) error { defer close(sp.quit) defer sp.close() - var kadDepth int + if r.doSync { + var kadDepth int - r.delivery.overlay.EachConn(nil, 256, func(addr network.OverlayConn, po int, nn bool) bool { - // TODO: stop or expose by kademlia - if nn { - kadDepth = po - } - return true - }) - - kad, ok := r.delivery.overlay.(*network.Kademlia) - if !ok { - return fmt.Errorf("Not a Kademlia!") - } - - var startPo int - var endPo int - var i int - var err error - - //iterate over each bin and solicit needed subscription to bins - kad.EachBin(r.addr.Over(), pot.DefaultPof(256), 0, func(po, size int, f func(func(val pot.Val, i int) bool) bool) bool { - - //identify begin and start index of the bin(s) we want to subscribe to - if po < kadDepth { - //not nn - endPo = po - if i > 0 { - startPo = endPo + 1 + r.delivery.overlay.EachConn(nil, 256, func(addr network.OverlayConn, po int, nn bool) bool { + // TODO: stop or expose by kademlia + if nn { + kadDepth = po } - } else if endPo < kadDepth || endPo == 0 { - if po == 0 && kadDepth == 0 { - startPo = endPo - } else { - startPo = endPo + 1 - } - endPo = maxPO + return true + }) + + kad, ok := r.delivery.overlay.(*network.Kademlia) + if !ok { + return fmt.Errorf("Not a Kademlia!") } - // now iterate and subscribe - for bin := po - startPo; bin <= endPo; bin++ { + var startPo int + var endPo int + var i int + var err error - f(func(val pot.Val, i int) bool { - // a := val.(network.OverlayPeer) - log.Debug(fmt.Sprintf("Requesting subscription by: registry %s from peer %s for bin: %d", r.addr.ID(), p.ID(), bin)) + //iterate over each bin and solicit needed subscription to bins + kad.EachBin(r.addr.Over(), pot.DefaultPof(256), 0, func(po, size int, f func(func(val pot.Val, i int) bool) bool) bool { - err = r.RequestSubscription(p.ID(), NewStream("SYNC", []byte{uint8(bin)}, true), &Range{}, Top) - if err != nil { - log.Error(fmt.Sprintf("Error in RequestSubsciption! %v", err)) - return false + //identify begin and start index of the bin(s) we want to subscribe to + if po < kadDepth { + //not nn + endPo = po + if i > 0 { + startPo = endPo + 1 } - return true - }) - } - i++ - return true - }) + } else if endPo < kadDepth || endPo == 0 { + if po == 0 && kadDepth == 0 { + startPo = endPo + } else { + startPo = endPo + 1 + } + endPo = maxPO + } + + // now iterate and subscribe + for bin := po - startPo; bin <= endPo; bin++ { + + f(func(val pot.Val, i int) bool { + // a := val.(network.OverlayPeer) + log.Debug(fmt.Sprintf("Requesting subscription by: registry %s from peer %s for bin: %d", r.addr.ID(), p.ID(), bin)) + + err = r.RequestSubscription(p.ID(), NewStream("SYNC", []byte{uint8(bin)}, true), &Range{}, Top) + if err != nil { + log.Error(fmt.Sprintf("Error in RequestSubsciption! %v", err)) + return false + } + return true + }) + } + i++ + return true + }) + } return sp.Run(sp.HandleMsg) } @@ -539,13 +542,11 @@ func (r *Registry) APIs() []rpc.API { } func (r *Registry) Start(server *p2p.Server) error { - r.api.dpa.Start() log.Info("Streamer started") return nil } func (r *Registry) Stop() error { - r.api.dpa.Stop() return nil } @@ -566,14 +567,11 @@ func getHistoryStream(s Stream) Stream { type API struct { streamer *Registry - dpa *storage.DPA } -func NewAPI(r *Registry, store storage.ChunkStore) *API { - dpa := storage.NewDPA(store, storage.NewChunkerParams()) +func NewAPI(r *Registry) *API { return &API{ streamer: r, - dpa: dpa, } } diff --git a/swarm/network/stream/syncer_test.go b/swarm/network/stream/syncer_test.go index 9afe38adbb..0a3ba57a31 100644 --- a/swarm/network/stream/syncer_test.go +++ b/swarm/network/stream/syncer_test.go @@ -163,6 +163,7 @@ func testSyncBetweenNodes(t *testing.T, nodes, conns, chunkCount int, skipCheck // start syncing, i.e., subscribe to upstream peers po 1 bin sid := sim.IDs[j+1] return client.CallContext(ctx, nil, "stream_subscribeStream", sid, NewStream("SYNC", []byte{1}, false), &Range{From: 0, To: 0}, Top) + return nil }) if err != nil { return err diff --git a/swarm/swarm.go b/swarm/swarm.go index d2c073c33c..2aa2cffb92 100644 --- a/swarm/swarm.go +++ b/swarm/swarm.go @@ -156,9 +156,7 @@ func NewSwarm(ctx *node.ServiceContext, backend chequebook.Backend, config *api. if err != nil { return } - self.streamer = stream.NewRegistry(addr, delivery, self.lstore, stateStore, false) - stream.RegisterSwarmSyncerServer(self.streamer, db) - stream.RegisterSwarmSyncerClient(self.streamer, db) + self.streamer = stream.NewRegistry(addr, delivery, db, stateStore, false, true) self.bzz = network.NewBzz(bzzconfig, to, stateStore, stream.Spec, self.streamer.Run)