From 459fafde7b52f9600dd4c2750372a6751b9b297d Mon Sep 17 00:00:00 2001 From: Fabio Barone Date: Wed, 11 Apr 2018 16:50:25 -0500 Subject: [PATCH 1/4] swarm: implemented retrieval test --- swarm/network/stream/common_test.go | 25 +- .../network/stream/snapshot_retrieval_test.go | 378 ++++++++++++++++++ swarm/network/stream/snapshot_sync_test.go | 78 ++-- 3 files changed, 442 insertions(+), 39 deletions(-) create mode 100644 swarm/network/stream/snapshot_retrieval_test.go diff --git a/swarm/network/stream/common_test.go b/swarm/network/stream/common_test.go index 8428fd3fb4..a67913774a 100644 --- a/swarm/network/stream/common_test.go +++ b/swarm/network/stream/common_test.go @@ -50,14 +50,20 @@ var ( peerCount func(discover.NodeID) int adapter = flag.String("adapter", "sim", "type of simulation: sim|socket|exec|docker") loglevel = flag.Int("loglevel", 2, "verbosity of logs") + nodes = flag.Int("nodes", 0, "number of nodes") + chunks = flag.Int("chunks", 0, "number of chunks") ) var ( - defaultSkipCheck bool - waitPeerErrC chan error - chunkSize = 4096 - registries map[discover.NodeID]*TestRegistry - createStoreFunc func(id discover.NodeID, addr *network.BzzAddr) (storage.ChunkStore, error) + defaultSkipCheck bool + defaultDoRetrieve bool + waitPeerErrC chan error + chunkSize = 4096 + registries map[discover.NodeID]*TestRegistry + createStoreFunc func(id discover.NodeID, addr *network.BzzAddr) (storage.ChunkStore, error) + getRetrieveFunc = defaultRetrieveFunc + doRetrieve = defaultDoRetrieve + subscriptionCount = 0 ) var services = adapters.Services{ @@ -90,19 +96,24 @@ func NewStreamerService(ctx *adapters.ServiceContext) (node.Service, error) { delivery := NewDelivery(kad, db) deliveries[id] = delivery r := NewRegistry(addr, delivery, db, state.NewMemStore(), &RegistryOptions{ - SkipCheck: defaultSkipCheck, + SkipCheck: defaultSkipCheck, + DoRetrieve: defaultDoRetrieve, }) RegisterSwarmSyncerServer(r, db) RegisterSwarmSyncerClient(r, db) go func() { waitPeerErrC <- waitForPeers(r, 1*time.Second, peerCount(id)) }() - dpa := storage.NewDPA(storage.NewNetStore(store, nil), storage.NewDPAParams()) + dpa := storage.NewDPA(storage.NewNetStore(store, getRetrieveFunc(id)), storage.NewDPAParams()) testRegistry := &TestRegistry{Registry: r, dpa: dpa} registries[id] = testRegistry return testRegistry, nil } +func defaultRetrieveFunc(id discover.NodeID) func(chunk *storage.Chunk) error { + return nil +} + func datadirsCleanup() { for _, id := range ids { os.RemoveAll(datadirs[id]) diff --git a/swarm/network/stream/snapshot_retrieval_test.go b/swarm/network/stream/snapshot_retrieval_test.go new file mode 100644 index 0000000000..6a932584ee --- /dev/null +++ b/swarm/network/stream/snapshot_retrieval_test.go @@ -0,0 +1,378 @@ +// 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 stream + +import ( + "context" + // crand "crypto/rand" + "fmt" + //"io" + "math/rand" + //"sync" + "testing" + "time" + + "github.com/ethereum/go-ethereum/log" + //"github.com/ethereum/go-ethereum/p2p" + "github.com/ethereum/go-ethereum/p2p/discover" + "github.com/ethereum/go-ethereum/p2p/simulations" + //"github.com/ethereum/go-ethereum/p2p/simulations/adapters" + //"github.com/ethereum/go-ethereum/pot" + //"github.com/ethereum/go-ethereum/rpc" + "github.com/ethereum/go-ethereum/swarm/network" + streamTesting "github.com/ethereum/go-ethereum/swarm/network/stream/testing" + "github.com/ethereum/go-ethereum/swarm/storage" +) + +func initRetrievalTest() { + toAddr = func(id discover.NodeID) *network.BzzAddr { + addr := network.NewAddrFromNodeID(id) + return addr + } + createStoreFunc = createTestLocalStorageForId + //local stores + stores = make(map[discover.NodeID]storage.ChunkStore) + //data directories for each node and store + datadirs = make(map[discover.NodeID]string) + //deliveries for each node + deliveries = make(map[discover.NodeID]*Delivery) + getRetrieveFunc = func(id discover.NodeID) func(chunk *storage.Chunk) error { + return func(chunk *storage.Chunk) error { + skipCheck := true + //fmt.Println(fmt.Sprintf("-- %s", id)) + return deliveries[id].RequestFromPeers(chunk.Key[:], skipCheck) + } + } + //registries, map of discover.NodeID to its streamer + registries = make(map[discover.NodeID]*TestRegistry) + //channel to wait for peers connected + //not needed for this test but required from common_test for NewStreamService + waitPeerErrC = make(chan error) + //also not needed for this test but required for NewStreamService + peerCount = func(id discover.NodeID) int { + if ids[0] == id || ids[len(ids)-1] == id { + return 1 + } + return 2 + } +} + +func TestRetrieval(t *testing.T) { + + if *nodes != 0 && *chunks != 0 { + retrievalTest(t, *chunks, *nodes) + } else { + var nodeCnt []int + var chnkCnt []int + if *longrunning { + nodeCnt = []int{16, 32, 128} + chnkCnt = []int{4, 32, 256} + } else { + nodeCnt = []int{16} + chnkCnt = []int{32} + } + for _, n := range nodeCnt { + for _, c := range chnkCnt { + retrievalTest(t, c, n) + } + } + } +} + +func retrievalTest(t *testing.T, chunkCount int, nodeCount int) { + //test live and NO history + log.Info("Testing live and no history", "chunkCount", chunkCount, "nodeCount", nodeCount) + live = true + history = false + err := runRetrievalTest(chunkCount, nodeCount) + if err != nil { + t.Fatal(err) + } + //test history only + log.Info("Testing history only", "chunkCount", chunkCount, "nodeCount", nodeCount) + live = false + history = true + err = runRetrievalTest(chunkCount, nodeCount) + if err != nil { + t.Fatal(err) + } + //finally test live and history + log.Info("Testing live and history", "chunkCount", chunkCount, "nodeCount", nodeCount) + live = true + err = runRetrievalTest(chunkCount, nodeCount) + if err != nil { + t.Fatal(err) + } +} + +/* +The test generates the given number of chunks, +then uploads these to a random node. +Afterwards for every chunk generated, the nearest node addresses +are identified, syncing is started, and finally 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 runRetrievalTest(chunkCount int, nodeCount int) error { + initRetrievalTest() + ids = make([]discover.NodeID, nodeCount) + disconnectC := make(chan error) + quitC := make(chan struct{}) + conf = &synctestConfig{} + //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) + conf.chunks = make([]storage.Key, 0) + //load nodes from the snapshot file + net, err := initNetWithSnapshot(nodeCount) + if err != nil { + return err + } + //do cleanup after test is terminated + defer func() { + doRetrieve = defaultDoRetrieve + //shutdown the snapshot network + net.Shutdown() + //after the test, clean up local stores initialized with createLocalStoreForId + localStoreCleanup() + //finally clear all data directories + datadirsCleanup() + }() + //get the nodes of the network + nodes := net.GetNodes() + //select one index at random... + idx := rand.Intn(len(nodes)) + //...and get the the node at that index + //this is the node selected for upload + uploadNode := nodes[idx] + //iterate over all nodes... + for c := 0; c < len(nodes); c++ { + //create an array of discovery nodeIDS + ids[c] = nodes[c].ID() + a := network.ToOverlayAddr(ids[c].Bytes()) + //append it to the array of all overlay addresses + conf.addrs = append(conf.addrs, a) + conf.idToAddrMap[ids[c]] = a + conf.addrToIdMap[string(a)] = ids[c] + } + + //needed for healthy call + 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{}) + + trigger := make(chan discover.NodeID) + action := func(ctx context.Context) error { + ticker := time.NewTicker(200 * time.Millisecond) + defer ticker.Stop() + for range ticker.C { + healthy := true + for _, id := range ids { + r := registries[id] + //PeerPot for this node + pp := ppmap[id] + //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)) + if !h.GotNN || !h.Full { + healthy = false + break + } + } + if healthy { + break + } + } + + if history { + log.Info("Uploading for history") + //If testing only history, we upload the chunk(s) first + conf.chunks, err = uploadFileToSingleNodeStore(uploadNode.ID(), chunkCount) + if err != nil { + return err + } + } + + //variables needed to wait for all subscriptions established before uploading + errc := make(chan error) + + //now setup and start event watching in order to know when we can upload + ctx, watchCancel := context.WithTimeout(context.Background(), MAX_TIMEOUT*time.Second) + defer watchCancel() + + log.Info("Setting up stream subscription") + // each node Subscribes to each other's swarmChunkServerStreamName + for j, id := range ids { + log.Trace(fmt.Sprintf("Subscribe to subscription events: %d", j)) + client, err := net.GetNode(id).Client() + if err != nil { + return err + } + //watch for peers disconnecting + err = streamTesting.WatchDisconnections(id, client, disconnectC, quitC) + if err != nil { + return err + } + + watchSubscriptionEvents(ctx, id, client, errc) + } + + for j, id := range ids { + log.Trace(fmt.Sprintf("Start syncing and stream subscriptions: %d", j)) + client, err := net.GetNode(id).Client() + if err != nil { + return err + } + //start syncing! + var cnt int + err = client.CallContext(ctx, &cnt, "stream_startSyncing") + if err != nil { + return err + } + subscriptionCount += cnt + for snid := range registries[id].peers { + subscriptionCount++ + err = client.CallContext(ctx, nil, "stream_subscribeStream", snid, NewStream(swarmChunkServerStreamName, "", false), nil, Top) + if err != nil { + return err + } + } + } + + //now wait until the number of expected subscriptions has been finished + for err := range errc { + if err != nil { + return err + } + subscriptionCount-- + if subscriptionCount == 0 { + break + } + } + + log.Info("Stream subscriptions successfully requested, action terminated") + + if live { + //now upload the chunks to the selected random single node + chnks, err := uploadFileToSingleNodeStore(uploadNode.ID(), chunkCount) + if err != nil { + return err + } + conf.chunks = append(conf.chunks, chnks...) + } + + return nil + } + + chunkSize := storage.DefaultChunkSize + + //check defines what will be checked during the test + check := func(ctx context.Context, id discover.NodeID) (bool, error) { + + if id == uploadNode.ID() { + return true, nil + } + + select { + case <-ctx.Done(): + return false, ctx.Err() + case e := <-disconnectC: + log.Error(e.Error()) + return false, fmt.Errorf("Disconnect event detected, network unhealthy") + default: + } + log.Trace(fmt.Sprintf("Checking node: %s", id)) + //if there are more than one chunk, test only succeeds if all expected chunks are found + allSuccess := true + + dpa := registries[id].dpa + for _, chnk := range conf.chunks { + reader := dpa.Retrieve(chnk) + if s, err := reader.Size(nil); err != nil || s != chunkSize { + allSuccess = false + log.Warn("Retrieve error", "err", err, "chunk", chnk, "nodeId", id) + } else { + log.Debug(fmt.Sprintf("Chunk %x found", chnk)) + } + } + return allSuccess, nil + } + + //for each tick, run the checks on all nodes + timingTicker := time.NewTicker(5 * time.Second) + defer timingTicker.Stop() + go func() { + for range timingTicker.C { + for i := 0; i < len(ids); i++ { + log.Trace(fmt.Sprintf("triggering step %d, id %s", i, ids[i])) + trigger <- ids[i] + } + } + }() + + log.Info("Starting simulation run...") + + timeout := MAX_TIMEOUT * time.Second + ctx, cancel := context.WithTimeout(context.Background(), timeout) + defer cancel() + + //run the simulation + result := simulations.NewSimulation(net).Run(ctx, &simulations.Step{ + Action: action, + Trigger: trigger, + Expect: &simulations.Expectation{ + Nodes: ids, + Check: check, + }, + }) + //close(quitC) + if result.Error != nil { + return result.Error + } + return nil +} + +//upload a file(chunks) +/* +func uploadRandomChunks(net *simulations.Network, chunkCount int) error { + log.Debug(fmt.Sprintf("Uploading to node id: %s", id)) + lstore := stores[id] + size := chunkCount * chunkSize + dpa := storage.NewDPA(lstore, storage.NewChunkerParams()) + dpa.Start() + rootHash, wait, err := dpa.Store(io.LimitReader(crand.Reader, int64(size)), int64(size)) + wait() + if err != nil { + return nil, err + } + + defer dpa.Stop() + + return rootHash, nil +} +*/ diff --git a/swarm/network/stream/snapshot_sync_test.go b/swarm/network/stream/snapshot_sync_test.go index 2df870cc77..7932aab48e 100644 --- a/swarm/network/stream/snapshot_sync_test.go +++ b/swarm/network/stream/snapshot_sync_test.go @@ -25,7 +25,6 @@ import ( "io/ioutil" "math/rand" "os" - "sync" "testing" "time" @@ -54,8 +53,6 @@ var ( datadirs map[discover.NodeID]string ppmap map[string]*network.PeerPot - globalWg sync.WaitGroup - live bool history bool @@ -107,20 +104,24 @@ func initSyncTest() { } } -//This file executes a number of tests with the syntax -//TestSyncing_x_y -//x is the number of chunks which will be uploaded -//y is the number of nodes for the test -func TestSyncing_4_32(t *testing.T) { testSyncing(t, 4, 32) } -func TestSyncing_32_16(t *testing.T) { testSyncing(t, 32, 16) } - -func TestLongRunningSyncing(t *testing.T) { - if *longrunning { - chnkCnt := []int{1, 8, 32, 256, 1024} - nCnt := []int{16, 32, 64, 128, 256} - +//This file executes a number of syncing tests +//node and chunk number can be provided via flags +func TestSyncing(t *testing.T) { + if *nodes != 0 && *chunks != 0 { + log.Info(fmt.Sprintf("Running test with %d chunks and %d nodes...", *chunks, *nodes)) + testSyncing(t, *chunks, *nodes) + } else { + var nodeCnt []int + var chnkCnt []int + if *longrunning { + chnkCnt = []int{1, 8, 32, 256, 1024} + nodeCnt = []int{16, 32, 64, 128, 256} + } else { + chnkCnt = []int{4, 32} + nodeCnt = []int{32, 16} + } for _, chnk := range chnkCnt { - for _, n := range nCnt { + for _, n := range nodeCnt { log.Info(fmt.Sprintf("Long running test with %d chunks and %d nodes...", chnk, n)) testSyncing(t, chnk, n) } @@ -291,7 +292,7 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { log.Info("Setting up stream subscription") // each node Subscribes to each other's swarmChunkServerStreamName for j, id := range ids { - log.Trace(fmt.Sprintf("subscribe: %d", j)) + log.Trace(fmt.Sprintf("Subscribe to subscription events: %d", j)) client, err := net.GetNode(id).Client() if err != nil { return err @@ -331,23 +332,34 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { <-wdDoneC rpcSubscriptionsWg.Done() }() - //start syncing! - err = client.CallContext(ctx, nil, "stream_startSyncing") + } + + for j, id := range ids { + log.Trace(fmt.Sprintf("Start syncing subscriptions: %d", j)) + client, err := net.GetNode(id).Client() if err != nil { return err } + //start syncing! + var cnt int + err = client.CallContext(ctx, &cnt, "stream_startSyncing") + if err != nil { + return err + } + subscriptionCount += cnt } //now wait until the number of expected subscriptions has been finished - go func() { - globalWg.Wait() - errc <- nil - }() - - err := <-errc - if err != nil { - return err + for err := range errc { + if err != nil { + return err + } + subscriptionCount-- + if subscriptionCount == 0 { + break + } } + log.Info("Stream subscriptions successfully requested") if live { //now upload the chunks to the selected random single node @@ -443,7 +455,7 @@ func (r *TestRegistry) GetKad(ctx context.Context) string { } //the server func to start syncing -func (r *TestRegistry) StartSyncing(ctx context.Context) error { +func (r *TestRegistry) StartSyncing(ctx context.Context) (int, error) { var err error if log.Lvl(*loglevel) == log.LvlDebug { @@ -459,9 +471,10 @@ func (r *TestRegistry) StartSyncing(ctx context.Context) error { kad, ok := r.delivery.overlay.(*network.Kademlia) if !ok { - return fmt.Errorf("Not a Kademlia!") + return 0, fmt.Errorf("Not a Kademlia!") } + subCnt := 0 //iterate over each bin and solicit needed subscription to bins kad.EachBin(r.addr.Over(), pof, 0, func(conn network.OverlayConn, po int) bool { //identify begin and start index of the bin(s) we want to subscribe to @@ -471,7 +484,7 @@ func (r *TestRegistry) StartSyncing(ctx context.Context) error { histRange = &Range{} } - globalWg.Add(1) + subCnt++ err = r.RequestSubscription(conf.addrToIdMap[string(conn.Address())], NewStream("SYNC", FormatSyncBinKey(uint8(po)), live), histRange, Top) if err != nil { log.Error(fmt.Sprintf("Error in RequestSubsciption! %v", err)) @@ -480,7 +493,7 @@ func (r *TestRegistry) StartSyncing(ctx context.Context) error { return true }) - return nil + return subCnt, nil } //map chunk keys to addresses which are responsible @@ -616,6 +629,7 @@ func watchSubscriptionEvents(ctx context.Context, id discover.NodeID, client *rp return } c := make(chan struct{}) + go func() { defer func() { log.Trace("watch subscription events: unsubscribe", "id", id) @@ -636,7 +650,7 @@ func watchSubscriptionEvents(ctx context.Context, id discover.NodeID, client *rp case e := <-events: //just catch SubscribeMsg if e.Type == p2p.PeerEventTypeMsgRecv && e.Protocol == "stream" && e.MsgCode != nil && *e.MsgCode == 4 { - globalWg.Done() + errc <- nil } case err := <-sub.Err(): if err != nil { From f5302571398cfbe8e57c7174407a2c6ec60add83 Mon Sep 17 00:00:00 2001 From: Fabio Barone Date: Thu, 12 Apr 2018 08:45:42 -0500 Subject: [PATCH 2/4] swarm: consolidate snapshot sync and retrieval tests --- .../network/stream/snapshot_retrieval_test.go | 104 ++++++++++-------- swarm/network/stream/snapshot_sync_test.go | 91 +++++++++------ 2 files changed, 115 insertions(+), 80 deletions(-) diff --git a/swarm/network/stream/snapshot_retrieval_test.go b/swarm/network/stream/snapshot_retrieval_test.go index 6a932584ee..b5bbc16651 100644 --- a/swarm/network/stream/snapshot_retrieval_test.go +++ b/swarm/network/stream/snapshot_retrieval_test.go @@ -17,31 +17,26 @@ package stream import ( "context" - // crand "crypto/rand" "fmt" - //"io" "math/rand" - //"sync" "testing" "time" "github.com/ethereum/go-ethereum/log" - //"github.com/ethereum/go-ethereum/p2p" "github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/p2p/simulations" - //"github.com/ethereum/go-ethereum/p2p/simulations/adapters" - //"github.com/ethereum/go-ethereum/pot" - //"github.com/ethereum/go-ethereum/rpc" "github.com/ethereum/go-ethereum/swarm/network" streamTesting "github.com/ethereum/go-ethereum/swarm/network/stream/testing" "github.com/ethereum/go-ethereum/swarm/storage" ) func initRetrievalTest() { + //global func to get overlay address from discover ID toAddr = func(id discover.NodeID) *network.BzzAddr { addr := network.NewAddrFromNodeID(id) return addr } + //global func to create local store createStoreFunc = createTestLocalStorageForId //local stores stores = make(map[discover.NodeID]storage.ChunkStore) @@ -49,16 +44,15 @@ func initRetrievalTest() { datadirs = make(map[discover.NodeID]string) //deliveries for each node deliveries = make(map[discover.NodeID]*Delivery) + //global retrieve func getRetrieveFunc = func(id discover.NodeID) func(chunk *storage.Chunk) error { return func(chunk *storage.Chunk) error { skipCheck := true - //fmt.Println(fmt.Sprintf("-- %s", id)) return deliveries[id].RequestFromPeers(chunk.Key[:], skipCheck) } } //registries, map of discover.NodeID to its streamer registries = make(map[discover.NodeID]*TestRegistry) - //channel to wait for peers connected //not needed for this test but required from common_test for NewStreamService waitPeerErrC = make(chan error) //also not needed for this test but required for NewStreamService @@ -70,17 +64,27 @@ func initRetrievalTest() { } } +//This test is a retrieval test for nodes. +//One node is randomly selected to be the pivot node. +//A configurable number of chunks and nodes can be +//provided to the test, the number of chunks is uploaded +//to the pivot node and other nodes try to retrieve the chunk(s). +//Number of chunks and nodes can be provided via commandline too. func TestRetrieval(t *testing.T) { - + //if nodes/chunks have been provided via commandline, + //run the tests with these values if *nodes != 0 && *chunks != 0 { retrievalTest(t, *chunks, *nodes) } else { var nodeCnt []int var chnkCnt []int + //if the `longrunning` flag has been provided + //run more test combinations if *longrunning { nodeCnt = []int{16, 32, 128} chnkCnt = []int{4, 32, 256} } else { + //default test nodeCnt = []int{16} chnkCnt = []int{32} } @@ -92,6 +96,7 @@ func TestRetrieval(t *testing.T) { } } +//Every test runs 3 times, a live, a history, and a live AND history func retrievalTest(t *testing.T, chunkCount int, nodeCount int) { //test live and NO history log.Info("Testing live and no history", "chunkCount", chunkCount, "nodeCount", nodeCount) @@ -119,22 +124,33 @@ func retrievalTest(t *testing.T, chunkCount int, nodeCount int) { } /* -The test generates the given number of chunks, -then uploads these to a random node. -Afterwards for every chunk generated, the nearest node addresses -are identified, syncing is started, and finally we verify -that the nodes closer to the chunk addresses actually do have -the chunks in their local stores. +The test generates the given number of chunks. + +The upload is done by dependency to the global +`live` and `history` variables; + +If `live` is set, first stream subscriptions are established, then +upload to a random node. + +If `history` is enabled, first upload then build up subscriptions. 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. +kademlia network. Nevertheless a health check runs in the +simulation's `action` function. + +The snapshot should have 'streamer' in its service list. */ func runRetrievalTest(chunkCount int, nodeCount int) error { + //for every run (live, history), int the variables initRetrievalTest() + //the ids of the snapshot nodes, initiate only now as we need nodeCount ids = make([]discover.NodeID, nodeCount) + //channel to check for disconnection errors disconnectC := make(chan error) + //channel to close disconnection watcher routine quitC := make(chan struct{}) + //the test conf (using same as in `snapshot_sync_test` conf = &synctestConfig{} //map of discover ID to indexes of chunks expected at that ID conf.idToChunksMap = make(map[discover.NodeID][]int) @@ -142,6 +158,7 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { conf.idToAddrMap = make(map[discover.NodeID][]byte) //map of overlay address to discover ID conf.addrToIdMap = make(map[string]discover.NodeID) + //array where the generated chunk hashes will be stored conf.chunks = make([]storage.Key, 0) //load nodes from the snapshot file net, err := initNetWithSnapshot(nodeCount) @@ -179,13 +196,11 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { //needed for healthy call 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{}) - trigger := make(chan discover.NodeID) + //simulation action action := func(ctx context.Context) error { + //first run the health check on all nodes, + //wait until nodes are all healthy ticker := time.NewTicker(200 * time.Millisecond) defer ticker.Stop() for range ticker.C { @@ -226,7 +241,14 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { defer watchCancel() log.Info("Setting up stream subscription") - // each node Subscribes to each other's swarmChunkServerStreamName + //We need two iterations, one to subscribe to the subscription events + //(so we know when setup phase is finished), and one to + //actually run the stream subscriptions. We can't do it in the same iteration, + //because while the first nodes in the loop are setting up subscriptions, + //the latter ones have not subscribed to listen to peer events yet, + //and then we miss events. + + //first iteration: setup disconnection watcher and subscribe to peer events for j, id := range ids { log.Trace(fmt.Sprintf("Subscribe to subscription events: %d", j)) client, err := net.GetNode(id).Client() @@ -239,9 +261,11 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { return err } + //check for `SubscribeMsg` events to know when setup phase is complete watchSubscriptionEvents(ctx, id, client, errc) } + //second iteration: start syncing and setup stream subscriptions for j, id := range ids { log.Trace(fmt.Sprintf("Start syncing and stream subscriptions: %d", j)) client, err := net.GetNode(id).Client() @@ -254,7 +278,10 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { 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 + //now also add the number of RETRIEVAL_REQUEST subscriptions for snid := range registries[id].peers { subscriptionCount++ err = client.CallContext(ctx, nil, "stream_subscribeStream", snid, NewStream(swarmChunkServerStreamName, "", false), nil, Top) @@ -265,11 +292,15 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { } //now wait until the number of expected subscriptions has been finished + //`watchSubscriptionEvents` will write with a `nil` value to errc + //every time a `SubscriptionMsg` has been received for err := range errc { if err != nil { return err } + //`nil` received, decrement count subscriptionCount-- + //all subscriptions received if subscriptionCount == 0 { break } @@ -294,6 +325,7 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { //check defines what will be checked during the test check := func(ctx context.Context, id discover.NodeID) (bool, error) { + //don't check the uploader node if id == uploadNode.ID() { return true, nil } @@ -310,9 +342,12 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { //if there are more than one chunk, test only succeeds if all expected chunks are found allSuccess := true + //check on the node's dpa (netstore) dpa := registries[id].dpa + //check all chunks for _, chnk := range conf.chunks { reader := dpa.Retrieve(chnk) + //assuming that reading the Size of the chunk is enough to know we found it if s, err := reader.Size(nil); err != nil || s != chunkSize { allSuccess = false log.Warn("Retrieve error", "err", err, "chunk", chnk, "nodeId", id) @@ -350,29 +385,10 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { Check: check, }, }) - //close(quitC) + if result.Error != nil { return result.Error } + return nil } - -//upload a file(chunks) -/* -func uploadRandomChunks(net *simulations.Network, chunkCount int) error { - log.Debug(fmt.Sprintf("Uploading to node id: %s", id)) - lstore := stores[id] - size := chunkCount * chunkSize - dpa := storage.NewDPA(lstore, storage.NewChunkerParams()) - dpa.Start() - rootHash, wait, err := dpa.Store(io.LimitReader(crand.Reader, int64(size)), int64(size)) - wait() - if err != nil { - return nil, err - } - - defer dpa.Stop() - - return rootHash, nil -} -*/ diff --git a/swarm/network/stream/snapshot_sync_test.go b/swarm/network/stream/snapshot_sync_test.go index 7932aab48e..62260f4895 100644 --- a/swarm/network/stream/snapshot_sync_test.go +++ b/swarm/network/stream/snapshot_sync_test.go @@ -82,7 +82,7 @@ func initSyncTest() { addr := network.NewAddrFromNodeID(id) return addr } - + //global func to create local store createStoreFunc = createTestLocalStorageForId //local stores stores = make(map[discover.NodeID]storage.ChunkStore) @@ -92,7 +92,6 @@ func initSyncTest() { deliveries = make(map[discover.NodeID]*Delivery) //registries, map of discover.NodeID to its streamer registries = make(map[discover.NodeID]*TestRegistry) - //channel to wait for peers connected //not needed for this test but required from common_test for NewStreamService waitPeerErrC = make(chan error) //also not needed for this test but required for NewStreamService @@ -104,19 +103,29 @@ func initSyncTest() { } } -//This file executes a number of syncing tests -//node and chunk number can be provided via flags +//This test is a syncing test for nodes. +//One node is randomly selected to be the pivot node. +//A configurable number of chunks and nodes can be +//provided to the test, the number of chunks is uploaded +//to the pivot node, and we check that nodes get the chunks +//they are expected to store based on the syncing protocol. +//Number of chunks and nodes can be provided via commandline too. func TestSyncing(t *testing.T) { + //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)) testSyncing(t, *chunks, *nodes) } 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{16, 32, 64, 128, 256} } else { + //default test chnkCnt = []int{4, 32} nodeCnt = []int{32, 16} } @@ -129,11 +138,9 @@ func TestSyncing(t *testing.T) { } } -//do run the tests +//Do run the tests +//Every test runs 3 times, a live, a history, and a live AND history func testSyncing(t *testing.T, chunkCount int, nodeCount int) { - initSyncTest() - ids = make([]discover.NodeID, nodeCount) - //test live and NO history log.Info("Testing live and no history") live = true @@ -160,12 +167,19 @@ func testSyncing(t *testing.T, chunkCount int, nodeCount int) { } /* -The test generates the given number of chunks, -then uploads these to a random node. -Afterwards for every chunk generated, the nearest node addresses -are identified, syncing is started, and finally we verify -that the nodes closer to the chunk addresses actually do have -the chunks in their local stores. +The test generates the given number of chunks + +The upload is done by dependency to the global +`live` and `history` variables; + +If `live` is set, first stream subscriptions are established, then +upload to a random node. + +If `history` is enabled, first upload then build up subscriptions. + +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 @@ -180,6 +194,9 @@ For every test run, a series of three tests will be executed: are uploaded twice, once before and once after subscriptions */ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { + initSyncTest() + //the ids of the snapshot nodes, initiate only now as we need nodeCount + ids = make([]discover.NodeID, nodeCount) //initialize the test struct conf = &synctestConfig{} //map of discover ID to indexes of chunks expected at that ID @@ -188,12 +205,13 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { conf.idToAddrMap = make(map[discover.NodeID][]byte) //map of overlay address to discover ID conf.addrToIdMap = make(map[string]discover.NodeID) + //array where the generated chunk hashes will be stored conf.chunks = make([]storage.Key, 0) - //First load the snapshot from the file + //channel to trigger node checks in the simulation trigger := make(chan discover.NodeID) - // channel to signal simulation initialisation with action call complete - // or node disconnections + //channel to check for disconnection errors disconnectC := make(chan error) + //channel to close disconnection watcher routine quitC := make(chan struct{}) //load nodes from the snapshot file @@ -246,6 +264,8 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { //define the action to be performed before the test checks: start syncing action := func(ctx context.Context) error { + //first run the health check on all nodes, + //wait until nodes are all healthy ticker := time.NewTicker(200 * time.Millisecond) defer ticker.Stop() for range ticker.C { @@ -290,7 +310,15 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { defer watchCancel() log.Info("Setting up stream subscription") - // each node Subscribes to each other's swarmChunkServerStreamName + + //We need two iterations, one to subscribe to the subscription events + //(so we know when setup phase is finished), and one to + //actually run the stream subscriptions. We can't do it in the same iteration, + //because while the first nodes in the loop are setting up subscriptions, + //the latter ones have not subscribed to listen to peer events yet, + //and then we miss events. + + //first iteration: setup disconnection watcher and subscribe to peer events for j, id := range ids { log.Trace(fmt.Sprintf("Subscribe to subscription events: %d", j)) client, err := net.GetNode(id).Client() @@ -309,19 +337,6 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { rpcSubscriptionsWg.Done() }() - if log.Lvl(*loglevel) >= log.LvlTrace { - //this will print the kademlia tables of all nodes - //to only print the kademlia of the pivot node, - //use: if j == idx {} - var kt string - err = client.CallContext(ctx, &kt, "stream_getKad") - if err != nil { - return err - } - - log.Debug("kad table " + node.ID().String()) - log.Debug(kt) - } //watch for peers disconnecting wdDoneC, err := streamTesting.WatchDisconnections(id, client, disconnectC, quitC) if err != nil { @@ -334,6 +349,7 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { }() } + //second iteration: start syncing for j, id := range ids { log.Trace(fmt.Sprintf("Start syncing subscriptions: %d", j)) client, err := net.GetNode(id).Client() @@ -346,15 +362,20 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { 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 } //now wait until the number of expected subscriptions has been finished + //`watchSubscriptionEvents` will write with a `nil` value to errc for err := range errc { if err != nil { return err } + //`nil` received, decrement count subscriptionCount-- + //all subscriptions received if subscriptionCount == 0 { break } @@ -449,12 +470,10 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { return nil } -//Show kademlia of uploading node for debugging -func (r *TestRegistry) GetKad(ctx context.Context) string { - return r.delivery.overlay.String() -} - //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 (r *TestRegistry) StartSyncing(ctx context.Context) (int, error) { var err error From e6f01d019d6b2f87ad867105a9aad9afc69513c6 Mon Sep 17 00:00:00 2001 From: Fabio Barone Date: Thu, 12 Apr 2018 21:47:19 -0500 Subject: [PATCH 3/4] swarm: added multiple file retrieval test --- swarm/network/stream/common_test.go | 4 +- .../network/stream/snapshot_retrieval_test.go | 384 +++++++++++++++++- swarm/network/stream/snapshot_sync_test.go | 32 +- 3 files changed, 387 insertions(+), 33 deletions(-) diff --git a/swarm/network/stream/common_test.go b/swarm/network/stream/common_test.go index a67913774a..73745f4ced 100644 --- a/swarm/network/stream/common_test.go +++ b/swarm/network/stream/common_test.go @@ -56,13 +56,11 @@ var ( var ( defaultSkipCheck bool - defaultDoRetrieve bool waitPeerErrC chan error chunkSize = 4096 registries map[discover.NodeID]*TestRegistry createStoreFunc func(id discover.NodeID, addr *network.BzzAddr) (storage.ChunkStore, error) getRetrieveFunc = defaultRetrieveFunc - doRetrieve = defaultDoRetrieve subscriptionCount = 0 ) @@ -97,7 +95,7 @@ func NewStreamerService(ctx *adapters.ServiceContext) (node.Service, error) { deliveries[id] = delivery r := NewRegistry(addr, delivery, db, state.NewMemStore(), &RegistryOptions{ SkipCheck: defaultSkipCheck, - DoRetrieve: defaultDoRetrieve, + DoRetrieve: false, }) RegisterSwarmSyncerServer(r, db) RegisterSwarmSyncerClient(r, db) diff --git a/swarm/network/stream/snapshot_retrieval_test.go b/swarm/network/stream/snapshot_retrieval_test.go index b5bbc16651..70a32f224e 100644 --- a/swarm/network/stream/snapshot_retrieval_test.go +++ b/swarm/network/stream/snapshot_retrieval_test.go @@ -19,6 +19,7 @@ import ( "context" "fmt" "math/rand" + "strings" "testing" "time" @@ -30,6 +31,13 @@ import ( "github.com/ethereum/go-ethereum/swarm/storage" ) +//constants for random file generation +const ( + minFileSize = 2 + maxFileSize = 40 + charset = ".,/.!@#$%^&*()-=:;<>abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789" +) + func initRetrievalTest() { //global func to get overlay address from discover ID toAddr = func(id discover.NodeID) *network.BzzAddr { @@ -64,6 +72,30 @@ func initRetrievalTest() { } } +//This test is a retrieval test for nodes. +//A configurable number of nodes can be +//provided to the test. +//Files are uploaded to nodes, other nodes try to retrieve the file +//Number of nodes can be provided via commandline too. +func TestFileRetrieval(t *testing.T) { + if *nodes != 0 { + fileRetrievalTest(t, *nodes) + } else { + var nodeCnt []int + //if the `longrunning` flag has been provided + //run more test combinations + if *longrunning { + nodeCnt = []int{16, 32, 128, 256} + } else { + //default test + nodeCnt = []int{16} + } + for _, n := range nodeCnt { + fileRetrievalTest(t, n) + } + } +} + //This test is a retrieval test for nodes. //One node is randomly selected to be the pivot node. //A configurable number of chunks and nodes can be @@ -96,6 +128,33 @@ func TestRetrieval(t *testing.T) { } } +//Every test runs 3 times, a live, a history, and a live AND history +func fileRetrievalTest(t *testing.T, nodeCount int) { + //test live and NO history + log.Info("Testing live and no history", "nodeCount", nodeCount) + live = true + history = false + err := runFileRetrievalTest(nodeCount) + if err != nil { + t.Fatal(err) + } + //test history only + log.Info("Testing history only", "nodeCount", nodeCount) + live = false + history = true + err = runFileRetrievalTest(nodeCount) + if err != nil { + t.Fatal(err) + } + //finally test live and history + log.Info("Testing live and history", "nodeCount", nodeCount) + live = true + err = runFileRetrievalTest(nodeCount) + if err != nil { + t.Fatal(err) + } +} + //Every test runs 3 times, a live, a history, and a live AND history func retrievalTest(t *testing.T, chunkCount int, nodeCount int) { //test live and NO history @@ -123,6 +182,273 @@ func retrievalTest(t *testing.T, chunkCount int, nodeCount int) { } } +/* + +The upload is done by dependency to the global +`live` and `history` variables; + +If `live` is set, first stream subscriptions are established, +then files are uploaded to nodes. + +If `history` is enabled, first upload files, then build up subscriptions. + +The test loads a snapshot file to construct the swarm network, +assuming that the snapshot file identifies a healthy +kademlia network. Nevertheless a health check runs in the +simulation's `action` function. + +The snapshot should have 'streamer' in its service list. +*/ +func runFileRetrievalTest(nodeCount int) error { + //for every run (live, history), int the variables + initRetrievalTest() + //the ids of the snapshot nodes, initiate only now as we need nodeCount + ids = make([]discover.NodeID, nodeCount) + //channel to check for disconnection errors + disconnectC := make(chan error) + //channel to close disconnection watcher routine + quitC := make(chan struct{}) + //the test conf (using same as in `snapshot_sync_test` + conf = &synctestConfig{} + //map of overlay address to discover ID + conf.addrToIdMap = make(map[string]discover.NodeID) + //array where the generated chunk hashes will be stored + conf.hashes = make([]storage.Key, 0) + //load nodes from the snapshot file + net, err := initNetWithSnapshot(nodeCount) + if err != nil { + return err + } + //do cleanup after test is terminated + defer func() { + //shutdown the snapshot network + net.Shutdown() + //after the test, clean up local stores initialized with createLocalStoreForId + localStoreCleanup() + //finally clear all data directories + datadirsCleanup() + }() + //get the nodes of the network + nodes := net.GetNodes() + //iterate over all nodes... + for c := 0; c < len(nodes); c++ { + //create an array of discovery nodeIDS + ids[c] = nodes[c].ID() + a := network.ToOverlayAddr(ids[c].Bytes()) + //append it to the array of all overlay addresses + conf.addrs = append(conf.addrs, a) + conf.addrToIdMap[string(a)] = ids[c] + } + + //needed for healthy call + ppmap = network.NewPeerPot(testMinProxBinSize, ids, conf.addrs) + + //an array for the random files + var randomFiles []string + //channel to signal when the upload has finished + uploadFinished := make(chan struct{}) + //channel to trigger new node checks + trigger := make(chan discover.NodeID) + //simulation action + action := func(ctx context.Context) error { + //first run the health check on all nodes, + //wait until nodes are all healthy + ticker := time.NewTicker(200 * time.Millisecond) + defer ticker.Stop() + for range ticker.C { + healthy := true + for _, id := range ids { + r := registries[id] + //PeerPot for this node + pp := ppmap[id] + //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)) + if !h.GotNN || !h.Full { + healthy = false + break + } + } + if healthy { + break + } + } + + if history { + log.Info("Uploading for history") + //If testing only history, we upload the chunk(s) first + conf.hashes, randomFiles, err = uploadFilesToNodes(nodes) + if err != nil { + return err + } + } + + //variables needed to wait for all subscriptions established before uploading + errc := make(chan error) + + //now setup and start event watching in order to know when we can upload + ctx, watchCancel := context.WithTimeout(context.Background(), MaxTimeout*time.Second) + defer watchCancel() + + log.Info("Setting up stream subscription") + //We need two iterations, one to subscribe to the subscription events + //(so we know when setup phase is finished), and one to + //actually run the stream subscriptions. We can't do it in the same iteration, + //because while the first nodes in the loop are setting up subscriptions, + //the latter ones have not subscribed to listen to peer events yet, + //and then we miss events. + + //first iteration: setup disconnection watcher and subscribe to peer events + for j, id := range ids { + log.Trace(fmt.Sprintf("Subscribe to subscription events: %d", j)) + client, err := net.GetNode(id).Client() + if err != nil { + return err + } + //watch for peers disconnecting + err = streamTesting.WatchDisconnections(id, client, disconnectC, quitC) + if err != nil { + return err + } + + //check for `SubscribeMsg` events to know when setup phase is complete + watchSubscriptionEvents(ctx, id, client, errc) + } + + //second iteration: start syncing and setup stream subscriptions + for j, id := range ids { + log.Trace(fmt.Sprintf("Start syncing and stream subscriptions: %d", j)) + client, err := net.GetNode(id).Client() + if err != nil { + return err + } + //start syncing! + var cnt int + err = client.CallContext(ctx, &cnt, "stream_startSyncing") + 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 + //now also add the number of RETRIEVAL_REQUEST subscriptions + for snid := range registries[id].peers { + subscriptionCount++ + err = client.CallContext(ctx, nil, "stream_subscribeStream", snid, NewStream(swarmChunkServerStreamName, "", false), nil, Top) + if err != nil { + return err + } + } + } + + //now wait until the number of expected subscriptions has been finished + //`watchSubscriptionEvents` will write with a `nil` value to errc + //every time a `SubscriptionMsg` has been received + for err := range errc { + if err != nil { + return err + } + //`nil` received, decrement count + subscriptionCount-- + //all subscriptions received + if subscriptionCount == 0 { + break + } + } + + log.Info("Stream subscriptions successfully requested, action terminated") + + if live { + //upload generated files to nodes + var hashes []storage.Key + var rfiles []string + hashes, rfiles, err = uploadFilesToNodes(nodes) + if err != nil { + return err + } + conf.hashes = append(conf.hashes, hashes...) + randomFiles = append(randomFiles, rfiles...) + //signal to the trigger loop that the upload has finished + uploadFinished <- struct{}{} + } + + return nil + } + + //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 e := <-disconnectC: + log.Error(e.Error()) + return false, fmt.Errorf("Disconnect event detected, network unhealthy") + default: + } + log.Trace(fmt.Sprintf("Checking node: %s", id)) + //if there are more than one chunk, test only succeeds if all expected chunks are found + allSuccess := true + + //check on the node's dpa (netstore) + dpa := registries[id].dpa + //check all chunks + for i, hash := range conf.hashes { + reader := dpa.Retrieve(hash) + //check that we can read the file size and that it corresponds to the generated file size + if s, err := reader.Size(nil); err != nil || s != int64(len(randomFiles[i])) { + allSuccess = false + log.Warn("Retrieve error", "err", err, "hash", hash, "nodeId", id) + } else { + log.Debug(fmt.Sprintf("File with root hash %x successfully retrieved", hash)) + } + } + + return allSuccess, nil + } + + //for each tick, run the checks on all nodes + timingTicker := time.NewTicker(5 * time.Second) + defer timingTicker.Stop() + go func() { + //for live upload, we should wait for uploads to have finished + //before starting to trigger the checks, due to file size + if live { + <-uploadFinished + } + for range timingTicker.C { + for i := 0; i < len(ids); i++ { + log.Trace(fmt.Sprintf("triggering step %d, id %s", i, ids[i])) + trigger <- ids[i] + } + } + }() + + log.Info("Starting simulation run...") + + timeout := MaxTimeout * time.Second + ctx, cancel := context.WithTimeout(context.Background(), timeout) + defer cancel() + + //run the simulation + result := simulations.NewSimulation(net).Run(ctx, &simulations.Step{ + Action: action, + Trigger: trigger, + Expect: &simulations.Expectation{ + Nodes: ids, + Check: check, + }, + }) + + if result.Error != nil { + return result.Error + } + + return nil +} + /* The test generates the given number of chunks. @@ -152,14 +478,10 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { quitC := make(chan struct{}) //the test conf (using same as in `snapshot_sync_test` conf = &synctestConfig{} - //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) //array where the generated chunk hashes will be stored - conf.chunks = make([]storage.Key, 0) + conf.hashes = make([]storage.Key, 0) //load nodes from the snapshot file net, err := initNetWithSnapshot(nodeCount) if err != nil { @@ -167,7 +489,6 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { } //do cleanup after test is terminated defer func() { - doRetrieve = defaultDoRetrieve //shutdown the snapshot network net.Shutdown() //after the test, clean up local stores initialized with createLocalStoreForId @@ -189,7 +510,6 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { a := network.ToOverlayAddr(ids[c].Bytes()) //append it to the array of all overlay addresses conf.addrs = append(conf.addrs, a) - conf.idToAddrMap[ids[c]] = a conf.addrToIdMap[string(a)] = ids[c] } @@ -227,7 +547,7 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { if history { log.Info("Uploading for history") //If testing only history, we upload the chunk(s) first - conf.chunks, err = uploadFileToSingleNodeStore(uploadNode.ID(), chunkCount) + conf.hashes, err = uploadFileToSingleNodeStore(uploadNode.ID(), chunkCount) if err != nil { return err } @@ -237,7 +557,7 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { errc := make(chan error) //now setup and start event watching in order to know when we can upload - ctx, watchCancel := context.WithTimeout(context.Background(), MAX_TIMEOUT*time.Second) + ctx, watchCancel := context.WithTimeout(context.Background(), MaxTimeout*time.Second) defer watchCancel() log.Info("Setting up stream subscription") @@ -314,7 +634,7 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { if err != nil { return err } - conf.chunks = append(conf.chunks, chnks...) + conf.hashes = append(conf.hashes, chnks...) } return nil @@ -345,7 +665,7 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { //check on the node's dpa (netstore) dpa := registries[id].dpa //check all chunks - for _, chnk := range conf.chunks { + for _, chnk := range conf.hashes { reader := dpa.Retrieve(chnk) //assuming that reading the Size of the chunk is enough to know we found it if s, err := reader.Size(nil); err != nil || s != chunkSize { @@ -372,7 +692,7 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { log.Info("Starting simulation run...") - timeout := MAX_TIMEOUT * time.Second + timeout := MaxTimeout * time.Second ctx, cancel := context.WithTimeout(context.Background(), timeout) defer cancel() @@ -392,3 +712,43 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { return nil } + +//upload generated files to nodes +//every node gets one file uploaded +func uploadFilesToNodes(nodes []*simulations.Node) ([]storage.Key, []string, error) { + nodeCnt := len(nodes) + log.Debug(fmt.Sprintf("Uploading %d files to nodes", nodeCnt)) + //array holding generated files + rfiles := make([]string, nodeCnt) + //array holding the root hashes of the files + rootkeys := make([]storage.Key, nodeCnt) + + //for every node, generate a file and upload + for i, n := range nodes { + id := n.ID() + dpa := registries[id].dpa + //generate a file + rfiles[i] = generateRandomFile() + //store it (upload it) on the dpa + rk, wait, err := dpa.Store(strings.NewReader(rfiles[i]), int64(len(rfiles[i])), false) + log.Debug("Uploaded random string file to node") + wait() + if err != nil { + return nil, nil, err + } + rootkeys[i] = rk + } + return rootkeys, rfiles, nil +} + +//generate a random file (string) +func generateRandomFile() string { + //generate a random file size between minFileSize and maxFileSize + fileSize := rand.Intn(maxFileSize-minFileSize) + minFileSize + log.Debug(fmt.Sprintf("Generated file with filesize %d kB", fileSize)) + b := make([]byte, fileSize*1024) + for i := range b { + b[i] = charset[rand.Intn(len(charset))] + } + return string(b) +} diff --git a/swarm/network/stream/snapshot_sync_test.go b/swarm/network/stream/snapshot_sync_test.go index 62260f4895..2fa55855a5 100644 --- a/swarm/network/stream/snapshot_sync_test.go +++ b/swarm/network/stream/snapshot_sync_test.go @@ -42,7 +42,7 @@ import ( ) const testMinProxBinSize = 2 -const MAX_TIMEOUT = 600 +const MaxTimeout = 600 var ( pof = pot.DefaultPof(256) @@ -61,10 +61,9 @@ var ( type synctestConfig struct { addrs [][]byte - chunks []storage.Key + hashes []storage.Key idToChunksMap map[discover.NodeID][]int chunksToNodesMap map[string][]int - idToAddrMap map[discover.NodeID][]byte addrToIdMap map[string]discover.NodeID } @@ -201,12 +200,10 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { conf = &synctestConfig{} //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) //array where the generated chunk hashes will be stored - conf.chunks = make([]storage.Key, 0) + conf.hashes = make([]storage.Key, 0) //channel to trigger node checks in the simulation trigger := make(chan discover.NodeID) //channel to check for disconnection errors @@ -254,7 +251,6 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { //the proximity calculation is on overlay addr, //the p2p/simulations check func triggers on discover.NodeID, //so we need to know which overlay addr maps to which nodeID - conf.idToAddrMap[ids[c]] = a conf.addrToIdMap[string(a)] = ids[c] } log.Info("Test config successfully initialized") @@ -297,7 +293,7 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { if err != nil { return err } - conf.chunks = append(conf.chunks, chunks...) + conf.hashes = append(conf.hashes, chunks...) //finally map chunks to the closest addresses mapKeysToNodes(conf) } @@ -306,7 +302,7 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { errc := make(chan error) //now setup and start event watching in order to know when we can upload - ctx, watchCancel := context.WithTimeout(context.Background(), MAX_TIMEOUT*time.Second) + ctx, watchCancel := context.WithTimeout(context.Background(), MaxTimeout*time.Second) defer watchCancel() log.Info("Setting up stream subscription") @@ -384,13 +380,13 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { log.Info("Stream subscriptions successfully requested") if live { //now upload the chunks to the selected random single node - chunks, err := uploadFileToSingleNodeStore(node.ID(), chunkCount) + hashes, err := uploadFileToSingleNodeStore(node.ID(), chunkCount) if err != nil { return err } - conf.chunks = append(conf.chunks, chunks...) + conf.hashes = append(conf.hashes, hashes...) //finally map chunks to the closest addresses - log.Debug(fmt.Sprintf("Uploaded chunks for live syncing: %v", conf.chunks)) + log.Debug(fmt.Sprintf("Uploaded chunks for live syncing: %v", conf.hashes)) mapKeysToNodes(conf) log.Info(fmt.Sprintf("Uploaded %d chunks to random single node", chunkCount)) } @@ -421,7 +417,7 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { //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] + chunk := conf.hashes[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 { @@ -449,7 +445,7 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { log.Info("Starting simulation run...") - timeout := MAX_TIMEOUT * time.Second + timeout := MaxTimeout * time.Second ctx, cancel := context.WithTimeout(context.Background(), timeout) defer cancel() @@ -527,11 +523,11 @@ func mapKeysToNodes(conf *synctestConfig) { np, _, _ = pot.Add(np, a, pof) } //for each address, run EachNeighbour on the chunk hashes pot to identify closest nodes - log.Trace(fmt.Sprintf("Generated hash chunk(s): %v", conf.chunks)) - for i := 0; i < len(conf.chunks); i++ { + log.Trace(fmt.Sprintf("Generated hash chunk(s): %v", conf.hashes)) + for i := 0; i < len(conf.hashes); i++ { pl := 256 //highest possible proximity var nns []int - np.EachNeighbour([]byte(conf.chunks[i]), pof, func(val pot.Val, po int) bool { + np.EachNeighbour([]byte(conf.hashes[i]), pof, func(val pot.Val, po int) bool { a := val.([]byte) if pl < 256 && pl != po { return false @@ -548,7 +544,7 @@ func mapKeysToNodes(conf *synctestConfig) { } return true }) - kmap[string(conf.chunks[i])] = nns + kmap[string(conf.hashes[i])] = nns } for addr, chunks := range nodemap { //this selects which chunks are expected to be found with the given node From b5db57448ebaf1f808c52ea1eddec81faa453765 Mon Sep 17 00:00:00 2001 From: Fabio Barone Date: Fri, 13 Apr 2018 10:02:24 -0500 Subject: [PATCH 4/4] swarm: solve merge conflicts for retrieval tests --- .../network/stream/snapshot_retrieval_test.go | 82 ++++++++++++++----- swarm/network/stream/snapshot_sync_test.go | 1 + 2 files changed, 62 insertions(+), 21 deletions(-) diff --git a/swarm/network/stream/snapshot_retrieval_test.go b/swarm/network/stream/snapshot_retrieval_test.go index 70a32f224e..2f9ae037d1 100644 --- a/swarm/network/stream/snapshot_retrieval_test.go +++ b/swarm/network/stream/snapshot_retrieval_test.go @@ -17,12 +17,15 @@ package stream import ( "context" + crand "crypto/rand" "fmt" "math/rand" "strings" + "sync" "testing" "time" + "github.com/ethereum/go-ethereum/common" "github.com/ethereum/go-ethereum/log" "github.com/ethereum/go-ethereum/p2p/discover" "github.com/ethereum/go-ethereum/p2p/simulations" @@ -35,7 +38,6 @@ import ( const ( minFileSize = 2 maxFileSize = 40 - charset = ".,/.!@#$%^&*()-=:;<>abcdefghijklmnopqrstuvwxyzABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789" ) func initRetrievalTest() { @@ -85,7 +87,7 @@ func TestFileRetrieval(t *testing.T) { //if the `longrunning` flag has been provided //run more test combinations if *longrunning { - nodeCnt = []int{16, 32, 128, 256} + nodeCnt = []int{16, 32, 128} } else { //default test nodeCnt = []int{16} @@ -219,6 +221,7 @@ func runFileRetrievalTest(nodeCount int) error { if err != nil { return err } + var rpcSubscriptionsWg sync.WaitGroup //do cleanup after test is terminated defer func() { //shutdown the snapshot network @@ -241,7 +244,7 @@ func runFileRetrievalTest(nodeCount int) error { } //needed for healthy call - ppmap = network.NewPeerPot(testMinProxBinSize, ids, conf.addrs) + ppmap = network.NewPeerPotMap(testMinProxBinSize, conf.addrs) //an array for the random files var randomFiles []string @@ -260,7 +263,8 @@ func runFileRetrievalTest(nodeCount int) error { for _, id := range ids { r := registries[id] //PeerPot for this node - pp := ppmap[id] + addr := common.Bytes2Hex(r.addr.OAddr) + pp := ppmap[addr] //call Healthy RPC h := r.delivery.overlay.Healthy(pp) //print info @@ -307,14 +311,27 @@ func runFileRetrievalTest(nodeCount int) error { if err != nil { return err } + wsDoneC := watchSubscriptionEvents(ctx, id, client, errc, quitC) + // doneC is nil, the error happened which is sent to errc channel, already + if wsDoneC == nil { + continue + } + rpcSubscriptionsWg.Add(1) + go func() { + <-wsDoneC + rpcSubscriptionsWg.Done() + }() + //watch for peers disconnecting - err = streamTesting.WatchDisconnections(id, client, disconnectC, quitC) + wdDoneC, err := streamTesting.WatchDisconnections(id, client, disconnectC, quitC) if err != nil { return err } - - //check for `SubscribeMsg` events to know when setup phase is complete - watchSubscriptionEvents(ctx, id, client, errc) + rpcSubscriptionsWg.Add(1) + go func() { + <-wdDoneC + rpcSubscriptionsWg.Done() + }() } //second iteration: start syncing and setup stream subscriptions @@ -396,7 +413,7 @@ func runFileRetrievalTest(nodeCount int) error { dpa := registries[id].dpa //check all chunks for i, hash := range conf.hashes { - reader := dpa.Retrieve(hash) + reader, _ := dpa.Retrieve(hash) //check that we can read the file size and that it corresponds to the generated file size if s, err := reader.Size(nil); err != nil || s != int64(len(randomFiles[i])) { allSuccess = false @@ -487,6 +504,7 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { if err != nil { return err } + var rpcSubscriptionsWg sync.WaitGroup //do cleanup after test is terminated defer func() { //shutdown the snapshot network @@ -514,7 +532,7 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { } //needed for healthy call - ppmap = network.NewPeerPot(testMinProxBinSize, ids, conf.addrs) + ppmap = network.NewPeerPotMap(testMinProxBinSize, conf.addrs) trigger := make(chan discover.NodeID) //simulation action @@ -528,7 +546,8 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { for _, id := range ids { r := registries[id] //PeerPot for this node - pp := ppmap[id] + addr := common.Bytes2Hex(network.ToOverlayAddr(id.Bytes())) + pp := ppmap[addr] //call Healthy RPC h := r.delivery.overlay.Healthy(pp) //print info @@ -575,14 +594,29 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { if err != nil { return err } + + //check for `SubscribeMsg` events to know when setup phase is complete + wsDoneC := watchSubscriptionEvents(ctx, id, client, errc, quitC) + // doneC is nil, the error happened which is sent to errc channel, already + if wsDoneC == nil { + continue + } + rpcSubscriptionsWg.Add(1) + go func() { + <-wsDoneC + rpcSubscriptionsWg.Done() + }() + //watch for peers disconnecting - err = streamTesting.WatchDisconnections(id, client, disconnectC, quitC) + wdDoneC, err := streamTesting.WatchDisconnections(id, client, disconnectC, quitC) if err != nil { return err } - - //check for `SubscribeMsg` events to know when setup phase is complete - watchSubscriptionEvents(ctx, id, client, errc) + rpcSubscriptionsWg.Add(1) + go func() { + <-wdDoneC + rpcSubscriptionsWg.Done() + }() } //second iteration: start syncing and setup stream subscriptions @@ -666,7 +700,7 @@ func runRetrievalTest(chunkCount int, nodeCount int) error { dpa := registries[id].dpa //check all chunks for _, chnk := range conf.hashes { - reader := dpa.Retrieve(chnk) + reader, _ := dpa.Retrieve(chnk) //assuming that reading the Size of the chunk is enough to know we found it if s, err := reader.Size(nil); err != nil || s != chunkSize { allSuccess = false @@ -723,12 +757,16 @@ func uploadFilesToNodes(nodes []*simulations.Node) ([]storage.Key, []string, err //array holding the root hashes of the files rootkeys := make([]storage.Key, nodeCnt) + var err error //for every node, generate a file and upload for i, n := range nodes { id := n.ID() dpa := registries[id].dpa //generate a file - rfiles[i] = generateRandomFile() + rfiles[i], err = generateRandomFile() + if err != nil { + return nil, nil, err + } //store it (upload it) on the dpa rk, wait, err := dpa.Store(strings.NewReader(rfiles[i]), int64(len(rfiles[i])), false) log.Debug("Uploaded random string file to node") @@ -742,13 +780,15 @@ func uploadFilesToNodes(nodes []*simulations.Node) ([]storage.Key, []string, err } //generate a random file (string) -func generateRandomFile() string { +func generateRandomFile() (string, error) { //generate a random file size between minFileSize and maxFileSize fileSize := rand.Intn(maxFileSize-minFileSize) + minFileSize log.Debug(fmt.Sprintf("Generated file with filesize %d kB", fileSize)) b := make([]byte, fileSize*1024) - for i := range b { - b[i] = charset[rand.Intn(len(charset))] + _, err := crand.Read(b) + if err != nil { + log.Error("Error generating random file.", "err", err) + return "", err } - return string(b) + return string(b), nil } diff --git a/swarm/network/stream/snapshot_sync_test.go b/swarm/network/stream/snapshot_sync_test.go index 2fa55855a5..ba0290ccd6 100644 --- a/swarm/network/stream/snapshot_sync_test.go +++ b/swarm/network/stream/snapshot_sync_test.go @@ -25,6 +25,7 @@ import ( "io/ioutil" "math/rand" "os" + "sync" "testing" "time"