swarm/network/stream: fix sync tests (#297)

This commit is contained in:
Janoš Guljaš 2018-03-05 18:42:47 +01:00 committed by GitHub
parent a3c93ca4c4
commit ad1ff36aaf
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
5 changed files with 79 additions and 69 deletions

View file

@ -32,6 +32,7 @@ import (
"github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/common"
"github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/log"
"github.com/ethereum/go-ethereum/node" "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/discover"
"github.com/ethereum/go-ethereum/p2p/simulations/adapters" "github.com/ethereum/go-ethereum/p2p/simulations/adapters"
p2ptest "github.com/ethereum/go-ethereum/p2p/testing" p2ptest "github.com/ethereum/go-ethereum/p2p/testing"
@ -78,13 +79,14 @@ func NewStreamerService(ctx *adapters.ServiceContext) (node.Service, error) {
db := storage.NewDBAPI(netStore) db := storage.NewDBAPI(netStore)
delivery := NewDelivery(kad, db) delivery := NewDelivery(kad, db)
deliveries[id] = delivery deliveries[id] = delivery
r := NewRegistry(addr, delivery, netStore, state.NewMemStore(), defaultSkipCheck) r := NewRegistry(addr, delivery, db, state.NewMemStore(), defaultSkipCheck, false)
RegisterSwarmSyncerServer(r, db) RegisterSwarmSyncerServer(r, db)
RegisterSwarmSyncerClient(r, db) RegisterSwarmSyncerClient(r, db)
go func() { go func() {
waitPeerErrC <- waitForPeers(r, 1*time.Second, peerCount(id)) 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) { 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) db := storage.NewDBAPI(netStore)
delivery := NewDelivery(to, db) delivery := NewDelivery(to, db)
streamer := NewRegistry(addr, delivery, localStore, state.NewMemStore(), defaultSkipCheck) streamer := NewRegistry(addr, delivery, db, state.NewMemStore(), defaultSkipCheck, false)
teardown := func() { teardown := func() {
streamer.Close() streamer.Close()
removeDataDir() removeDataDir()
@ -169,6 +171,7 @@ func (rrs *roundRobinStore) Close() {
type TestRegistry struct { type TestRegistry struct {
*Registry *Registry
dpa *storage.DPA
} }
func (r *TestRegistry) APIs() []rpc.API { 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) { 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 { type TestExternalRegistry struct {

View file

@ -52,7 +52,7 @@ func newIntervalsStreamerService(ctx *adapters.ServiceContext) (node.Service, er
db := storage.NewDBAPI(netStore) db := storage.NewDBAPI(netStore)
delivery := NewDelivery(kad, db) delivery := NewDelivery(kad, db)
deliveries[id] = delivery 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) { r.RegisterClientFunc(externalStreamName, func(p *Peer, t []byte, live bool) (Client, error) {
return newTestExternalClient(t, db), nil return newTestExternalClient(t, db), nil

View file

@ -56,23 +56,23 @@ type Registry struct {
clientFuncs map[string]func(*Peer, []byte, bool) (Client, error) clientFuncs map[string]func(*Peer, []byte, bool) (Client, error)
peers map[discover.NodeID]*Peer peers map[discover.NodeID]*Peer
delivery *Delivery delivery *Delivery
store storage.ChunkStore
intervalsStore state.Store intervalsStore state.Store
doSync bool
} }
// NewRegistry is Streamer constructor // 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{ streamer := &Registry{
addr: addr, addr: addr,
skipCheck: skipCheck, skipCheck: skipCheck,
store: store,
serverFuncs: make(map[string]func(*Peer, []byte, bool) (Server, error)), serverFuncs: make(map[string]func(*Peer, []byte, bool) (Server, error)),
clientFuncs: make(map[string]func(*Peer, []byte, bool) (Client, error)), clientFuncs: make(map[string]func(*Peer, []byte, bool) (Client, error)),
peers: make(map[discover.NodeID]*Peer), peers: make(map[discover.NodeID]*Peer),
delivery: delivery, delivery: delivery,
intervalsStore: intervalsStore, intervalsStore: intervalsStore,
doSync: doSync,
} }
streamer.api = NewAPI(streamer, streamer.store) streamer.api = NewAPI(streamer)
delivery.getPeer = streamer.getPeer delivery.getPeer = streamer.getPeer
streamer.RegisterServerFunc(swarmChunkServerStreamName, func(_ *Peer, _ []byte, _ bool) (Server, error) { streamer.RegisterServerFunc(swarmChunkServerStreamName, func(_ *Peer, _ []byte, _ bool) (Server, error) {
return NewSwarmChunkServer(delivery.db), nil 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) { streamer.RegisterClientFunc(swarmChunkServerStreamName, func(p *Peer, _ []byte, _ bool) (Client, error) {
return NewSwarmSyncerClient(p, delivery.db, nil) return NewSwarmSyncerClient(p, delivery.db, nil)
}) })
RegisterSwarmSyncerServer(streamer, db)
RegisterSwarmSyncerClient(streamer, db)
return streamer return streamer
} }
@ -214,7 +216,6 @@ func (r *Registry) PeerInfo(id discover.NodeID) interface{} {
} }
func (r *Registry) Close() error { func (r *Registry) Close() error {
r.store.Close()
return r.intervalsStore.Close() return r.intervalsStore.Close()
} }
@ -252,63 +253,65 @@ func (r *Registry) Run(p *network.BzzPeer) error {
defer close(sp.quit) defer close(sp.quit)
defer sp.close() 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 { r.delivery.overlay.EachConn(nil, 256, func(addr network.OverlayConn, po int, nn bool) bool {
// TODO: stop or expose by kademlia // TODO: stop or expose by kademlia
if nn { if nn {
kadDepth = po 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
} }
} else if endPo < kadDepth || endPo == 0 { return true
if po == 0 && kadDepth == 0 { })
startPo = endPo
} else { kad, ok := r.delivery.overlay.(*network.Kademlia)
startPo = endPo + 1 if !ok {
} return fmt.Errorf("Not a Kademlia!")
endPo = maxPO
} }
// now iterate and subscribe var startPo int
for bin := po - startPo; bin <= endPo; bin++ { var endPo int
var i int
var err error
f(func(val pot.Val, i int) bool { //iterate over each bin and solicit needed subscription to bins
// a := val.(network.OverlayPeer) kad.EachBin(r.addr.Over(), pot.DefaultPof(256), 0, func(po, size int, f func(func(val pot.Val, i int) bool) bool) bool {
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) //identify begin and start index of the bin(s) we want to subscribe to
if err != nil { if po < kadDepth {
log.Error(fmt.Sprintf("Error in RequestSubsciption! %v", err)) //not nn
return false endPo = po
if i > 0 {
startPo = endPo + 1
} }
return true } else if endPo < kadDepth || endPo == 0 {
}) if po == 0 && kadDepth == 0 {
} startPo = endPo
i++ } else {
return true 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) return sp.Run(sp.HandleMsg)
} }
@ -539,13 +542,11 @@ func (r *Registry) APIs() []rpc.API {
} }
func (r *Registry) Start(server *p2p.Server) error { func (r *Registry) Start(server *p2p.Server) error {
r.api.dpa.Start()
log.Info("Streamer started") log.Info("Streamer started")
return nil return nil
} }
func (r *Registry) Stop() error { func (r *Registry) Stop() error {
r.api.dpa.Stop()
return nil return nil
} }
@ -566,14 +567,11 @@ func getHistoryStream(s Stream) Stream {
type API struct { type API struct {
streamer *Registry streamer *Registry
dpa *storage.DPA
} }
func NewAPI(r *Registry, store storage.ChunkStore) *API { func NewAPI(r *Registry) *API {
dpa := storage.NewDPA(store, storage.NewChunkerParams())
return &API{ return &API{
streamer: r, streamer: r,
dpa: dpa,
} }
} }

View file

@ -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 // start syncing, i.e., subscribe to upstream peers po 1 bin
sid := sim.IDs[j+1] 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 client.CallContext(ctx, nil, "stream_subscribeStream", sid, NewStream("SYNC", []byte{1}, false), &Range{From: 0, To: 0}, Top)
return nil
}) })
if err != nil { if err != nil {
return err return err

View file

@ -156,9 +156,7 @@ func NewSwarm(ctx *node.ServiceContext, backend chequebook.Backend, config *api.
if err != nil { if err != nil {
return return
} }
self.streamer = stream.NewRegistry(addr, delivery, self.lstore, stateStore, false) self.streamer = stream.NewRegistry(addr, delivery, db, stateStore, false, true)
stream.RegisterSwarmSyncerServer(self.streamer, db)
stream.RegisterSwarmSyncerClient(self.streamer, db)
self.bzz = network.NewBzz(bzzconfig, to, stateStore, stream.Spec, self.streamer.Run) self.bzz = network.NewBzz(bzzconfig, to, stateStore, stream.Spec, self.streamer.Run)