From 5ca0180e8eb5b5d08a354cc485ca070a5d0c0729 Mon Sep 17 00:00:00 2001 From: Fabio Barone Date: Thu, 1 Mar 2018 21:04:28 -0500 Subject: [PATCH] swarm: Subscribe to bins with RequestSubscription --- swarm/network/stream/common_test.go | 3 +- swarm/network/stream/messages.go | 6 +- swarm/network/stream/peer.go | 5 + swarm/network/stream/snapshot_sync_test.go | 340 ++++++++++++--------- swarm/network/stream/streamer_test.go | 2 +- 5 files changed, 212 insertions(+), 144 deletions(-) diff --git a/swarm/network/stream/common_test.go b/swarm/network/stream/common_test.go index 420bb8c61f..c5b09a590a 100644 --- a/swarm/network/stream/common_test.go +++ b/swarm/network/stream/common_test.go @@ -114,11 +114,12 @@ func createTestLocalStorageForId(id discover.NodeID, addr *network.BzzAddr) (sto //local stores need to be cleaned up after the sim is done func localStoreCleanup() { - fmt.Println("Local store cleanup") + log.Info("Cleaning up...") for i := 0; i < len(ids); i++ { stores[ids[i]].Close() os.RemoveAll(datadirs[ids[i]]) } + log.Info("Local store cleanup done") } func newStreamerTester(t *testing.T) (*p2ptest.ProtocolTester, *Registry, *storage.LocalStore, func(), error) { diff --git a/swarm/network/stream/messages.go b/swarm/network/stream/messages.go index b58df047b5..9f2fb9553c 100644 --- a/swarm/network/stream/messages.go +++ b/swarm/network/stream/messages.go @@ -75,8 +75,8 @@ func (p *Peer) handleRequestSubscription(req *RequestSubscriptionMsg) (err error } func (p *Peer) handleRequestSubscription(req *RequestSubscriptionMsg) (err error) { - log.Debug(fmt.Sprintf("handleRequestSubscription: streamer %s to subscribe to %s", p.streamer.addr.ID(), p.ID())) - err = p.streamer.Subscribe(p.ID(), req.Stream, &Range{}, req.Priority) + log.Debug(fmt.Sprintf("handleRequestSubscription: streamer %s to subscribe to %s with stream %s", p.streamer.addr.ID(), p.ID(), req.Stream)) + err = p.streamer.Subscribe(p.ID(), req.Stream, req.History, req.Priority) if err != nil { return err } @@ -94,7 +94,7 @@ func (p *Peer) handleSubscribeMsg(req *SubscribeMsg) (err error) { } }() - log.Debug("received subscription", "peer", p.ID(), "stream", req.Stream, "history", req.History) + log.Debug("%s received subscription", "from", p.streamer.addr.ID(), "peer", p.ID(), "stream", req.Stream, "history", req.History) f, err := p.streamer.GetServerFunc(req.Stream.Name) if err != nil { diff --git a/swarm/network/stream/peer.go b/swarm/network/stream/peer.go index e41194ecaf..3b0a90b0f9 100644 --- a/swarm/network/stream/peer.go +++ b/swarm/network/stream/peer.go @@ -294,8 +294,13 @@ func (p *Peer) setClientParams(s Stream, params *clientParams) error { if p.clients[s] != nil { return fmt.Errorf("client %s already exists", s) } +<<<<<<< HEAD if p.clientParams[s] != nil { return fmt.Errorf("client params %s already set", s) +======= + if p.clientParams[sk] != nil { + return fmt.Errorf("client params %v already set, %s to %s", sk, p.streamer.addr.ID(), p.ID()) +>>>>>>> e13194f15... swarm: Subscribe to bins with RequestSubscription } p.clientParams[s] = params return nil diff --git a/swarm/network/stream/snapshot_sync_test.go b/swarm/network/stream/snapshot_sync_test.go index 164270a2d2..d402304e13 100644 --- a/swarm/network/stream/snapshot_sync_test.go +++ b/swarm/network/stream/snapshot_sync_test.go @@ -33,6 +33,7 @@ import ( "github.com/ethereum/go-ethereum/p2p/simulations/adapters" "github.com/ethereum/go-ethereum/pot" "github.com/ethereum/go-ethereum/swarm/network" + streamTesting "github.com/ethereum/go-ethereum/swarm/network/stream/testing" "github.com/ethereum/go-ethereum/swarm/storage" ) @@ -52,14 +53,15 @@ type synctestConfig struct { addrs [][]byte chunks []storage.Key retrievalMap map[string]map[string]time.Duration - nodesToChunksMap map[string][]int + idToChunksMap map[discover.NodeID][]int chunksToNodesMap map[string][]int idToAddrMap map[discover.NodeID][]byte addrToIdMap map[string]discover.NodeID } func init() { - rand.Seed(time.Now().Unix()) + //rand.Seed(time.Now().Unix()) + rand.Seed(100) initSyncTest() } @@ -175,14 +177,21 @@ 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. + +This tests LIVE syncing, as the file is uploaded *after* sync streams have been setup. +For HISTORY syncing a different test is needed. */ func runSyncTest(chunkCount int, nodeCount int) error { + //initialize the test struct conf = &synctestConfig{} //mapping of nearest node addresses for chunk hashes - //nodesToChunksMap = make(map[discover.NodeID][]storage.Key) conf.retrievalMap = make(map[string]map[string]time.Duration) + //map of discover ID to indexes of chunks expected at that ID + conf.idToChunksMap = make(map[discover.NodeID][]int) + //map of discover ID to kademlia overlay address conf.idToAddrMap = make(map[discover.NodeID][]byte) + //map of overlay address to discover ID conf.addrToIdMap = make(map[string]discover.NodeID) //First load the snapshot from the file net, err := initNetWithSnapshot(nodeCount) @@ -198,13 +207,15 @@ func runSyncTest(chunkCount int, nodeCount int) error { //...and get the the node at that index //this is the node selected for upload node := nodes[idx] + + log.Info("Initializing test config") //iterate over all nodes... for c := 0; c < len(nodes); c++ { - //create an array of discovery nodeIDS + //create an array of discovery node IDs ids[c] = nodes[c].ID() - //and a correspondent array of overlay addresses, - //later used for chunk proximity calculation + //get the kademlia overlay address from this ID a := network.ToOverlayAddr(ids[c].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 discover.NodeID, @@ -212,112 +223,142 @@ func runSyncTest(chunkCount int, nodeCount int) error { conf.idToAddrMap[ids[c]] = a conf.addrToIdMap[string(a)] = ids[c] } + log.Info("Test config successfully initialized") ppmap = network.NewPeerPot(testMinProxBinSize, ids, conf.addrs) // channel to signal simulation initialisation with action call complete // or node disconnections - //disconnectC := make(chan error) - //quitC := make(chan struct{}) + disconnectC := make(chan error) + quitC := make(chan struct{}) //after the test, clean up local stores initialized with createLocalStoreForId defer localStoreCleanup() - trigger := make(chan discover.NodeID) - //triggerCheck defines what will be checked during the test - triggerCheck := func(ctx context.Context, id discover.NodeID) (bool, error) { - select { - case <-ctx.Done(): - return false, ctx.Err() - //case <-disconnectC: - // log.Error("Disconnect event detected") - // return false, ctx.Err() - default: - } - - log.Debug(fmt.Sprintf("Checking node: %s", id)) - //select the local store for the given node - lstore := stores[id] - //if there are more than one chunk, test only succeeds if all expected chunks are found - allSuccess := true - //this selects which chunks are expected to be found with the given node - //localChunks := nodesToChunksMap[id] - localChunks := conf.nodesToChunksMap[string(conf.idToAddrMap[id])] - - //for each expected chunk, check if it is in the local store - for i := 0; i < len(localChunks); i++ { - //ignore zero chunks - chunk := conf.chunks[localChunks[i]] - if storage.IsZeroKey(chunk) { - continue - } - log.Debug(fmt.Sprintf("node has chunk: %s:", chunk)) - if _, err := lstore.Get(chunk); err != nil { - log.Error(fmt.Sprintf("Chunk %s NOT found for id %s", chunk, id)) - allSuccess = false - } else { - fmt.Println("^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^") - log.Info(fmt.Sprintf("Chunk %s FOUND for id %s", chunk, id)) - } - } - - return allSuccess, nil - } - - timeout := 300 * time.Second - ctx, cancel := context.WithTimeout(context.Background(), timeout) - defer cancel() //define the action to be performed before the test checks: start syncing action := func(ctx context.Context) error { // need to wait till an aynchronous process registers the peers in streamer.peers // that is used by Subscribe // the global peerCount function tells how many connections each node has // TODO: this is to be reimplemented with peerEvent watcher without global var - i := 0 - for err := range waitPeerErrC { - if err != nil { - return fmt.Errorf("error waiting for peers: %s", err) - } - i++ - if i == len(ids)-1 { - break - } - } - time.Sleep(10 * time.Second) + //TODO: VALIDATE THE ASSUMPTION THAT THE FOLLOWING CODE IS NOT NEEDED, + //AS THE SNAPSHOT CONSTRUCTS ALL CONNECTIONS DURING LOAD, SO WE DON'T NEED TO WAIT HERE? + + /* + i := 0 + for err := range waitPeerErrC { + fmt.Println("aaaa") + if err != nil { + return fmt.Errorf("error waiting for peers: %s", err) + } + i++ + if i == len(ids)-1 { + break + } + } + + // wait for connections + time.Sleep(5 * time.Second) + */ + + log.Info("Setting up stream subscription") + // each node Subscribes to each other's swarmChunkServerStreamName - for j := 0; j < len(ids); j++ { - log.Debug(fmt.Sprintf("subscribe: %d", j)) + for j, id := range ids { + log.Trace(fmt.Sprintf("subscribe: %d", j)) ctx, cancel := context.WithTimeout(ctx, 1*time.Second) defer cancel() - client, err := net.GetNode(ids[j]).Client() + client, err := net.GetNode(id).Client() + if err != nil { + return err + } + + if log.Lvl(*loglevel) == log.LvlDebug { + //print uploading node kademlia + if j == idx { + var kt string + err := client.CallContext(ctx, &kt, "stream_getKad") + if err != nil { + return err + } + + log.Debug("uploading node kad") + log.Debug(kt) + } + } + err = streamTesting.WatchDisconnections(id, client, disconnectC, quitC) if err != nil { return err } err = client.CallContext(ctx, nil, "stream_startSyncing") if err != nil { - log.Error(fmt.Sprintf("FAILED CallContext %v", err)) - return nil + return err } } + + log.Info("Stream subscriptions successfully requested") + // wait for subscritpions + //TODO: Implement a proper sync mechanism so that we don't need to Sleep() time.Sleep(10 * time.Second) //now upload the chunks to the selected random single node conf.chunks, err = uploadFileToSingleNodeStore(node.ID(), chunkCount) if err != nil { return err } + log.Info(fmt.Sprintf("Uploaded %d chunks to random single node", chunkCount)) //finally map chunks to the closest addresses - conf = mapKeysToNodes(conf) - log.Debug(fmt.Sprintf("%v", conf.nodesToChunksMap)) + mapKeysToNodes(conf) return nil } + trigger := make(chan discover.NodeID) + //check defines what will be checked during the test + check := func(ctx context.Context, id discover.NodeID) (bool, error) { + select { + case <-ctx.Done(): + return false, ctx.Err() + case <-disconnectC: + log.Error("Disconnect event detected") + return false, ctx.Err() + default: + } + + log.Trace(fmt.Sprintf("Checking node: %s", id)) + //select the local store for the given node + lstore := stores[id] + //if there are more than one chunk, test only succeeds if all expected chunks are found + allSuccess := true + + //all the chunk indexes which are supposed to be found for this node + localChunks := conf.idToChunksMap[id] + //for each expected chunk, check if it is in the local store + for _, ch := range localChunks { + //get the real chunk by the index in the index array + chunk := conf.chunks[ch] + log.Trace(fmt.Sprintf("node has chunk: %s:", chunk)) + //check if the expected chunk is indeed in the localstore + if _, err := lstore.Get(chunk); err != nil { + log.Debug(fmt.Sprintf("Chunk %s NOT found for id %s", chunk, id)) + allSuccess = false + } else { + log.Trace(fmt.Sprintf("Chunk %s FOUND for id %s", chunk, id)) + } + } + + return allSuccess, nil + } + + timeout := 120 * time.Second + ctx, cancel := context.WithTimeout(context.Background(), timeout) + defer cancel() + //for each tick, run the checks on all nodes go func() { ticker := time.NewTicker(time.Second * 1) for range ticker.C { for i := 0; i < len(ids); i++ { - log.Debug(fmt.Sprintf("triggering step %d, id %s", i, ids[i])) + log.Trace(fmt.Sprintf("triggering step %d, id %s", i, ids[i])) trigger <- ids[i] } } @@ -332,88 +373,98 @@ func runSyncTest(chunkCount int, nodeCount int) error { }() */ + log.Info("Starting simulation run...") //run the simulation result := simulations.NewSimulation(net).Run(ctx, &simulations.Step{ Action: action, Trigger: trigger, Expect: &simulations.Expectation{ Nodes: ids, - Check: triggerCheck, + Check: check, }, }) - //close(quitC) + close(quitC) if result.Error != nil { return result.Error } + log.Info("Simulation terminated") return nil } +func (r *TestRegistry) GetKad(ctx context.Context) string { + return r.delivery.overlay.String() +} + func (r *TestRegistry) StartSyncing(ctx context.Context) error { var err error - add := r.addr.ID() - pp := ppmap[add] - h := r.delivery.overlay.Healthy(pp) - fmt.Println("----------------------------------") - fmt.Println(r.delivery.overlay.String()) - fmt.Println(fmt.Sprintf("IS HEALTHY: %t", h.GotNN && h.KnowNN && h.Full)) + if log.Lvl(*loglevel) == log.LvlDebug { + //address of registry + add := r.addr.ID() + //PeerPot for this node + pp := ppmap[add] + //call Healthy RPC + h := r.delivery.overlay.Healthy(pp) + //print info + log.Debug(r.delivery.overlay.String()) + log.Debug(fmt.Sprintf("IS HEALTHY: %t", h.GotNN && h.KnowNN && h.Full)) + } - pos := make(map[int]discover.NodeID) + var kadDepth int r.delivery.overlay.EachConn(nil, 256, func(addr network.OverlayConn, po int, nn bool) bool { - lastPO := po + // TODO: stop or expose by kademlia if nn { - lastPO = maxPO - } - peerId := conf.addrToIdMap[string(addr.Address())] - fmt.Println(fmt.Sprintf("node %s has conn with %s at po %d and is nn: %t", r.addr.ID(), peerId, po, nn)) - pos[po] = peerId - for i := po; i <= lastPO; i++ { - err = r.Subscribe(peerId, NewStream("SYNC", []byte{byte(i)}, false), &Range{From: 0, To: 0}, Top) - if err != nil { - log.Error(fmt.Sprintf("Error subscribing! %v", err)) - return false - } + kadDepth = po } return true }) - prev := 0 + kad, ok := r.delivery.overlay.(*network.Kademlia) if !ok { return fmt.Errorf("Not a Kademlia!") } + var startPo int + var endPo int + var i int + + //iterate over each bin and solicit needed subscription to bins kad.EachBin(r.addr.Over(), pof, 0, func(po, size int, f func(func(val pot.Val, i int) bool) bool) bool { - skip := po - prev - /* - fmt.Println(prev) - fmt.Println(po) - fmt.Println(skip) - */ - remember := make(map[int]bool) - if skip > 1 { + + //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 { + 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 { - //for c := po + 1; c < po+skip; c++ { - for c := po - 1; c > po-skip; c-- { - //fmt.Println(c) - if exists, _ := remember[c]; exists { - continue - } - a := val.(network.OverlayPeer) - log.Warn(fmt.Sprintf("Request subscription for bin: %d", c)) - log.Debug(fmt.Sprintf("Requesting subscription by: registry %s from peer %s", r.addr.ID(), conf.addrToIdMap[string(a.Address())])) - err = r.RequestSubscription(conf.addrToIdMap[string(a.Address())], NewStream("SYNC", []byte{byte(uint8(c))}, false), Top) - if err != nil { - log.Error(fmt.Sprintf("Error subscribing! %v", err)) - return false - } - remember[c] = true + a := val.(network.OverlayPeer) + log.Debug(fmt.Sprintf("Requesting subscription by: registry %s from peer %s for bin: %d", r.addr.ID(), conf.addrToIdMap[string(a.Address())], bin)) + + err = r.RequestSubscription(conf.addrToIdMap[string(a.Address())], NewStream("SYNC", []byte{uint8(bin)}, true), &Range{}, Top) + if err != nil { + log.Error(fmt.Sprintf("Error in RequestSubsciption! %v", err)) + return false } return true }) } - prev = po + i++ return true }) @@ -439,36 +490,41 @@ func checkChunkIsAtNode(conf *synctestConfig) { log.Info("All chunks arrived at destination") for ch, n := range conf.retrievalMap { for a, t := range n { - log.Info(fmt.Sprintf("Chunk %s at node %s took %d ms", string(ch), string(a), t.Seconds()*1e3)) + log.Info(fmt.Sprintf("Chunk %s at node %s took %v ms", string(ch), string(a), t.Seconds()*1e3)) } } } } //map chunk keys to addresses which are responsible -func mapKeysToNodes(conf *synctestConfig) *synctestConfig { +func mapKeysToNodes(conf *synctestConfig) { kmap := make(map[string][]int) nodemap := make(map[string][]int) //build a pot for chunk hashes np := pot.NewPot(nil, 0) - mm := make(map[string]int) + indexmap := make(map[string]int) for i, a := range conf.addrs { - mm[string(a)] = i + indexmap[string(a)] = i np, _, _ = pot.Add(np, a, pof) } //for each address, run EachNeighbour on the chunk hashes pot to identify closest nodes - fmt.Println(conf.chunks) + log.Trace(fmt.Sprintf("Generated hash chunk(s): %v", conf.chunks)) for i := 0; i < len(conf.chunks); i++ { - pl := 256 //highest proximity + pl := 256 //highest possible proximity var nns []int np.EachNeighbour([]byte(conf.chunks[i]), pof, func(val pot.Val, po int) bool { a := val.([]byte) + if pl < 256 && pl != po { + return false + } if pl == 256 || pl == po { - fmt.Println(fmt.Sprintf("appending %s", conf.addrToIdMap[string(a)])) - nns = append(nns, mm[string(a)]) + log.Trace(fmt.Sprintf("appending %s", conf.addrToIdMap[string(a)])) + nns = append(nns, indexmap[string(a)]) nodemap[string(a)] = append(nodemap[string(a)], i) } if pl == 256 && len(nns) >= testMinProxBinSize { + //maxProxBinSize has been reached at this po, so save it + //we will add all other nodes at the same po pl = po } return true @@ -476,24 +532,27 @@ func mapKeysToNodes(conf *synctestConfig) *synctestConfig { kmap[conf.chunks[i].String()] = nns //log.Debug(fmt.Sprintf("Length for id %s: %d",ids[i],len(kmap[ids[i]]))) } - for k, v := range nodemap { - fmt.Print(fmt.Sprintf("Node %s: ", conf.addrToIdMap[k])) - for _, vv := range v { - fmt.Println(conf.chunks[vv]) + if log.Lvl(*loglevel) == log.LvlTrace { + for k, v := range nodemap { + log.Trace(fmt.Sprintf("Node %s: ", conf.addrToIdMap[k])) + for _, vv := range v { + log.Trace(fmt.Sprintf("%v", conf.chunks[vv])) + } + log.Trace(fmt.Sprintf("%v", conf.addrToIdMap[k])) } - fmt.Println(conf.addrToIdMap[k]) - fmt.Println("-------------------------------") - } - for k, v := range kmap { - fmt.Print(fmt.Sprintf("Chunk %s: ", k)) - for _, vv := range v { - fmt.Println(conf.addrToIdMap[string(conf.addrs[vv])]) + for k, v := range kmap { + log.Trace(fmt.Sprintf("Chunk %s: ", k)) + for _, vv := range v { + log.Trace(fmt.Sprintf("%v", conf.addrToIdMap[string(conf.addrs[vv])])) + } } - fmt.Println("###############################") } - conf.nodesToChunksMap = nodemap + for addr, chunks := range nodemap { + //this selects which chunks are expected to be found with the given node + conf.idToChunksMap[conf.addrToIdMap[addr]] = chunks + } + log.Debug(fmt.Sprintf("Map of expected chunks by ID: %v", conf.idToChunksMap)) conf.chunksToNodesMap = kmap - return conf } //upload a file(chunks) to a single local node store @@ -541,6 +600,8 @@ func initNetWithSnapshot(nodeCount int) (*simulations.Network, error) { a = adapters.NewSimAdapter(services) } + log.Info("Setting up Snapshot network") + net := simulations.NewNetwork(a, &simulations.NetworkConfig{ ID: "0", DefaultService: "streamer", @@ -564,5 +625,6 @@ func initNetWithSnapshot(nodeCount int) (*simulations.Network, error) { if err != nil { return nil, err } + log.Info("Snapshot loaded") return net, nil } diff --git a/swarm/network/stream/streamer_test.go b/swarm/network/stream/streamer_test.go index 2961a7e896..30d005f7f6 100644 --- a/swarm/network/stream/streamer_test.go +++ b/swarm/network/stream/streamer_test.go @@ -47,7 +47,7 @@ func TestStreamerRequestSubscription(t *testing.T) { } stream := NewStream("foo", nil, false) - err = streamer.RequestSubscription(tester.IDs[0], stream, Top) + err = streamer.RequestSubscription(tester.IDs[0], stream, &Range{}, Top) if err == nil || err.Error() != "stream foo not registered" { t.Fatalf("Expected error %v, got %v", "stream foo not registered", err) }