swarm/network/stream: address PR comments

This commit is contained in:
Fabio Barone 2018-12-26 16:28:32 -05:00
parent c32d961b0d
commit 689aad22fd
3 changed files with 22 additions and 284 deletions

View file

@ -106,43 +106,6 @@ func TestSyncingViaGlobalSync(t *testing.T) {
} }
} }
func TestSyncingViaDirectSubscribe(t *testing.T) {
if runtime.GOOS == "darwin" && os.Getenv("TRAVIS") == "true" {
t.Skip("Flaky on mac on travis")
}
//if nodes/chunks have been provided via commandline,
//run the tests with these values
if *nodes != 0 && *chunks != 0 {
log.Info(fmt.Sprintf("Running test with %d chunks and %d nodes...", *chunks, *nodes))
err := testSyncingViaDirectSubscribe(t, *chunks, *nodes)
if err != nil {
t.Fatal(err)
}
} else {
var nodeCnt []int
var chnkCnt []int
//if the `longrunning` flag has been provided
//run more test combinations
if *longrunning {
chnkCnt = []int{1, 8, 32, 256, 1024}
nodeCnt = []int{32, 16}
} else {
//default test
chnkCnt = []int{4, 32}
nodeCnt = []int{32, 16}
}
for _, chnk := range chnkCnt {
for _, n := range nodeCnt {
log.Info(fmt.Sprintf("Long running test with %d chunks and %d nodes...", chnk, n))
err := testSyncingViaDirectSubscribe(t, chnk, n)
if err != nil {
t.Fatal(err)
}
}
}
}
}
var simServiceMap = map[string]simulation.ServiceFunc{ var simServiceMap = map[string]simulation.ServiceFunc{
"streamer": streamerFunc, "streamer": streamerFunc,
} }
@ -323,235 +286,6 @@ func runSim(conf *synctestConfig, ctx context.Context, sim *simulation.Simulatio
}) })
} }
/*
The test generates the given number of chunks
For every chunk generated, the nearest node addresses
are identified, we verify that the nodes closer to the
chunk addresses actually do have the chunks in their local stores.
The test loads a snapshot file to construct the swarm network,
assuming that the snapshot file identifies a healthy
kademlia network. The snapshot should have 'streamer' in its service list.
*/
func testSyncingViaDirectSubscribe(t *testing.T, chunkCount int, nodeCount int) error {
sim := simulation.New(map[string]simulation.ServiceFunc{
"streamer": func(ctx *adapters.ServiceContext, bucket *sync.Map) (s node.Service, cleanup func(), err error) {
n := ctx.Config.Node()
addr := network.NewAddr(n)
store, datadir, err := createTestLocalStorageForID(n.ID(), addr)
if err != nil {
return nil, nil, err
}
bucket.Store(bucketKeyStore, store)
localStore := store.(*storage.LocalStore)
netStore, err := storage.NewNetStore(localStore, nil)
if err != nil {
return nil, nil, err
}
kad := network.NewKademlia(addr.Over(), network.NewKadParams())
delivery := NewDelivery(kad, netStore)
netStore.NewNetFetcherFunc = network.NewFetcherFactory(dummyRequestFromPeers, true).New
r := NewRegistry(addr.ID(), delivery, netStore, state.NewInmemoryStore(), &RegistryOptions{
Retrieval: RetrievalDisabled,
Syncing: SyncingRegisterOnly,
}, nil)
bucket.Store(bucketKeyRegistry, r)
fileStore := storage.NewFileStore(netStore, storage.NewFileStoreParams())
bucket.Store(bucketKeyFileStore, fileStore)
cleanup = func() {
os.RemoveAll(datadir)
netStore.Close()
r.Close()
}
return r, cleanup, nil
},
})
defer sim.Close()
ctx, cancelSimRun := context.WithTimeout(context.Background(), 2*time.Minute)
defer cancelSimRun()
conf := &synctestConfig{}
//map of discover ID to indexes of chunks expected at that ID
conf.idToChunksMap = make(map[enode.ID][]int)
//map of overlay address to discover ID
conf.addrToIDMap = make(map[string]enode.ID)
//array where the generated chunk hashes will be stored
conf.hashes = make([]storage.Address, 0)
err := sim.UploadSnapshot(fmt.Sprintf("testing/snapshot_%d.json", nodeCount))
if err != nil {
return err
}
if _, err := sim.WaitTillHealthy(ctx); err != nil {
return err
}
disconnections := sim.PeerEvents(
context.Background(),
sim.NodeIDs(),
simulation.NewPeerEventsFilter().Drop(),
)
var disconnected atomic.Value
go func() {
for d := range disconnections {
if d.Error != nil {
log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID)
disconnected.Store(true)
}
}
}()
result := sim.Run(ctx, func(ctx context.Context, sim *simulation.Simulation) error {
nodeIDs := sim.UpNodeIDs()
for _, n := range nodeIDs {
//get the kademlia overlay address from this ID
a := n.Bytes()
//append it to the array of all overlay addresses
conf.addrs = append(conf.addrs, a)
//the proximity calculation is on overlay addr,
//the p2p/simulations check func triggers on enode.ID,
//so we need to know which overlay addr maps to which nodeID
conf.addrToIDMap[string(a)] = n
}
var subscriptionCount int
filter := simulation.NewPeerEventsFilter().ReceivedMessages().Protocol("stream").MsgCode(4)
eventC := sim.PeerEvents(ctx, nodeIDs, filter)
for j, node := range nodeIDs {
log.Trace(fmt.Sprintf("Start syncing subscriptions: %d", j))
//start syncing!
item, ok := sim.NodeItem(node, bucketKeyRegistry)
if !ok {
return fmt.Errorf("No registry")
}
registry := item.(*Registry)
var cnt int
cnt, err = startSyncing(registry, conf)
if err != nil {
return err
}
//increment the number of subscriptions we need to wait for
//by the count returned from startSyncing (SYNC subscriptions)
subscriptionCount += cnt
}
for e := range eventC {
if e.Error != nil {
return e.Error
}
subscriptionCount--
if subscriptionCount == 0 {
break
}
}
//select a random node for upload
node := sim.Net.GetRandomUpNode()
item, ok := sim.NodeItem(node.ID(), bucketKeyStore)
if !ok {
return fmt.Errorf("No localstore")
}
lstore := item.(*storage.LocalStore)
hashes, err := uploadFileToSingleNodeStore(node.ID(), chunkCount, lstore)
if err != nil {
return err
}
conf.hashes = append(conf.hashes, hashes...)
mapKeysToNodes(conf)
if _, err := sim.WaitTillHealthy(ctx); err != nil {
return err
}
var globalStore mock.GlobalStorer
if *useMockStore {
globalStore = mockmem.NewGlobalStore()
}
// File retrieval check is repeated until all uploaded files are retrieved from all nodes
// or until the timeout is reached.
REPEAT:
for {
for _, id := range nodeIDs {
//for each expected chunk, check if it is in the local store
localChunks := conf.idToChunksMap[id]
for _, ch := range localChunks {
//get the real chunk by the index in the index array
chunk := conf.hashes[ch]
log.Trace(fmt.Sprintf("node has chunk: %s:", chunk))
//check if the expected chunk is indeed in the localstore
var err error
if *useMockStore {
//use the globalStore if the mockStore should be used; in that case,
//the complete localStore stack is bypassed for getting the chunk
_, err = globalStore.Get(common.BytesToAddress(id.Bytes()), chunk)
} else {
//use the actual localstore
item, ok := sim.NodeItem(id, bucketKeyStore)
if !ok {
return fmt.Errorf("Error accessing localstore")
}
lstore := item.(*storage.LocalStore)
_, err = lstore.Get(ctx, chunk)
}
if err != nil {
log.Debug(fmt.Sprintf("Chunk %s NOT found for id %s", chunk, id))
// Do not get crazy with logging the warn message
time.Sleep(500 * time.Millisecond)
continue REPEAT
}
log.Debug(fmt.Sprintf("Chunk %s IS FOUND for id %s", chunk, id))
}
}
return nil
}
})
if result.Error != nil {
return result.Error
}
if yes, ok := disconnected.Load().(bool); ok && yes {
t.Fatal("disconnect events received")
}
log.Info("Simulation ended")
return nil
}
//the server func to start syncing
//issues `RequestSubscriptionMsg` to peers, based on po, by iterating over
//the kademlia's `EachBin` function.
//returns the number of subscriptions requested
func startSyncing(r *Registry, conf *synctestConfig) (int, error) {
var err error
kad := r.delivery.kad
subCnt := 0
//iterate over each bin and solicit needed subscription to bins
kad.EachConn(nil, 255, func(p *network.Peer, po int, _ bool) bool {
//identify begin and start index of the bin(s) we want to subscribe to
subCnt++
err = r.RequestSubscription(conf.addrToIDMap[string(p.Address())], NewStream("SYNC", FormatSyncBinKey(uint8(po)), true), NewRange(0, 0), High)
if err != nil {
log.Error(fmt.Sprintf("Error in RequestSubsciption! %v", err))
return false
}
return true
})
return subCnt, nil
}
//map chunk keys to addresses which are responsible //map chunk keys to addresses which are responsible
func mapKeysToNodes(conf *synctestConfig) { func mapKeysToNodes(conf *synctestConfig) {
nodemap := make(map[string][]int) nodemap := make(map[string][]int)

View file

@ -87,9 +87,13 @@ type Registry struct {
intervalsStore state.Store intervalsStore state.Store
autoRetrieval bool // automatically subscribe to retrieve request stream autoRetrieval bool // automatically subscribe to retrieve request stream
maxPeerServers int maxPeerServers int
balance protocols.Balance // implements protocols.Balance, for accounting spec *protocols.Spec //this protocol's spec
prices protocols.Prices // implements protocols.Prices, provides prices to accounting balance protocols.Balance //implements protocols.Balance, for accounting
spec *protocols.Spec // this protocol's spec prices protocols.Prices //implements protocols.Prices, provides prices to accounting
// the subscriptionFunc is used to determine what to do in order to perform subscriptions
// usually we would start to really subscribe to nodes, but for tests other functionality may be needed
// (see TestRequestPeerSubscriptions in streamer_test.go)
subscriptionFunc func(p *network.Peer, bin uint8, subs map[enode.ID]map[Stream]struct{}) bool
} }
// RegistryOptions holds optional values for NewRegistry constructor. // RegistryOptions holds optional values for NewRegistry constructor.
@ -124,6 +128,8 @@ func NewRegistry(localID enode.ID, delivery *Delivery, syncChunkStore storage.Sy
maxPeerServers: options.MaxPeerServers, maxPeerServers: options.MaxPeerServers,
balance: balance, balance: balance,
} }
//assign the default subscription func: actually do request subscriptions from nodes
streamer.subscriptionFunc = streamer.doRequestSubscription
streamer.setupSpec() streamer.setupSpec()
streamer.api = NewAPI(streamer) streamer.api = NewAPI(streamer)
@ -467,7 +473,7 @@ func (r *Registry) updateSyncing() {
r.peersMu.RUnlock() r.peersMu.RUnlock()
// start requesting subscriptions from peers // start requesting subscriptions from peers
r.requestPeerSubscriptions(kad, subs, r.doRequestSubscription) r.requestPeerSubscriptions(kad, subs)
// remove SYNC servers that do not need to be subscribed // remove SYNC servers that do not need to be subscribed
for id, streams := range subs { for id, streams := range subs {
@ -498,10 +504,7 @@ func (r *Registry) updateSyncing() {
// * a map of subscriptions // * a map of subscriptions
// * the actual function to subscribe // * the actual function to subscribe
// (in case of the test, it doesn't do real subscriptions) // (in case of the test, it doesn't do real subscriptions)
func (r *Registry) requestPeerSubscriptions( func (r *Registry) requestPeerSubscriptions(kad *network.Kademlia, subs map[enode.ID]map[Stream]struct{}) {
kad *network.Kademlia,
subs map[enode.ID]map[Stream]struct{},
subscriptionFunc func(p *network.Peer, bin uint8, subs map[enode.ID]map[Stream]struct{}) bool) {
var startPo int var startPo int
var endPo int var endPo int
@ -527,7 +530,7 @@ func (r *Registry) requestPeerSubscriptions(
for bin := startPo; bin <= endPo; bin++ { for bin := startPo; bin <= endPo; bin++ {
//do the actual subscription //do the actual subscription
ok = subscriptionFunc(p, uint8(bin), subs) ok = r.subscriptionFunc(p, uint8(bin), subs)
} }
return ok return ok
}) })
@ -537,7 +540,7 @@ func (r *Registry) requestPeerSubscriptions(
func (r *Registry) doRequestSubscription(p *network.Peer, bin uint8, subs map[enode.ID]map[Stream]struct{}) bool { func (r *Registry) doRequestSubscription(p *network.Peer, bin uint8, subs map[enode.ID]map[Stream]struct{}) bool {
log.Debug(fmt.Sprintf("Requesting subscription by: registry %s from peer %s for bin: %d", r.addr, p.ID(), bin)) log.Debug(fmt.Sprintf("Requesting subscription by: registry %s from peer %s for bin: %d", r.addr, p.ID(), bin))
// bin is always less then 256 and it is safe to convert it to type uint8 // bin is always less then 256 and it is safe to convert it to type uint8
stream := NewStream("SYNC", FormatSyncBinKey(uint8(bin)), true) stream := NewStream("SYNC", FormatSyncBinKey(bin), true)
if streams, ok := subs[p.ID()]; ok { if streams, ok := subs[p.ID()]; ok {
// delete live and history streams from the map, so that it won't be removed with a Quit request // delete live and history streams from the map, so that it won't be removed with a Quit request
delete(streams, stream) delete(streams, stream)

View file

@ -962,7 +962,7 @@ func TestHasPriceImplementation(t *testing.T) {
TestRequestPeerSubscriptions is a unit test for stream's pull sync subscriptions. TestRequestPeerSubscriptions is a unit test for stream's pull sync subscriptions.
The test does: The test does:
* assign each known connected peer to a bin map * assign each connected peer to a bin map
* build up a known kademlia in advance * build up a known kademlia in advance
* run the EachConn function, which returns supposed subscription bins * run the EachConn function, which returns supposed subscription bins
* store all supposed bins per peer in a map * store all supposed bins per peer in a map
@ -1067,7 +1067,8 @@ func TestRequestPeerSubscriptions(t *testing.T) {
// create just a simple Registry object in order to be able to call... // create just a simple Registry object in order to be able to call...
r := &Registry{} r := &Registry{}
// ...the requestPeerSubscriptions function, which contains the logic for subscriptions // ...the requestPeerSubscriptions function, which contains the logic for subscriptions
r.requestPeerSubscriptions(k, nil, requestSubscriptionFunc) r.subscriptionFunc = requestSubscriptionFunc
r.requestPeerSubscriptions(k, nil)
// calculate the kademlia depth // calculate the kademlia depth
kdepth := k.NeighbourhoodDepth() kdepth := k.NeighbourhoodDepth()
@ -1077,27 +1078,27 @@ func TestRequestPeerSubscriptions(t *testing.T) {
// for every peer... // for every peer...
for _, peer := range peers { for _, peer := range peers {
// ...get its (fake) subscriptions // ...get its (fake) subscriptions
fakeSubs := fakeSubscriptions[peer] fakeSubsForPeer := fakeSubscriptions[peer]
// if the peer's bin is below the kademlia depth... // if the peer's bin is shallower than the kademlia depth...
if bin < kdepth { if bin < kdepth {
// (iterate all (fake) subscriptions) // (iterate all (fake) subscriptions)
for _, subbin := range fakeSubs { for _, subbin := range fakeSubsForPeer {
// ...only the peer's bin should be "subscribed" // ...only the peer's bin should be "subscribed"
// (and thus have only one subscription) // (and thus have only one subscription)
if subbin != bin || len(fakeSubs) != 1 { if subbin != bin || len(fakeSubsForPeer) != 1 {
t.Fatalf("Did not get expected subscription for bin < depth; bin of peer %s: %d, subscription: %d", peer, bin, subbin) t.Fatalf("Did not get expected subscription for bin < depth; bin of peer %s: %d, subscription: %d", peer, bin, subbin)
} }
} }
} else { //if the peer's bin is equal or higher than the kademlia depth... } else { //if the peer's bin is equal or higher than the kademlia depth...
// (iterate all (fake) subscriptions) // (iterate all (fake) subscriptions)
for i, subbin := range fakeSubs { for i, subbin := range fakeSubsForPeer {
// ...each bin from the peer's bin number up to k.MaxProxDisplay should be "subscribed" // ...each bin from the peer's bin number up to k.MaxProxDisplay should be "subscribed"
// as we start from depth we can use the iteration index to check // as we start from depth we can use the iteration index to check
if subbin != i+kdepth { if subbin != i+kdepth {
t.Fatalf("Did not get expected subscription for bin > depth; bin of peer %s: %d, subscription: %d", peer, bin, subbin) t.Fatalf("Did not get expected subscription for bin > depth; bin of peer %s: %d, subscription: %d", peer, bin, subbin)
} }
// the last "subscription" should be k.MaxProxDisplay // the last "subscription" should be k.MaxProxDisplay
if i == len(fakeSubs)-1 && subbin != k.MaxProxDisplay { if i == len(fakeSubsForPeer)-1 && subbin != k.MaxProxDisplay {
t.Fatalf("Expected last subscription to be: %d, but is: %d", k.MaxProxDisplay, subbin) t.Fatalf("Expected last subscription to be: %d, but is: %d", k.MaxProxDisplay, subbin)
} }
} }