From edcb0e1a30b0ad2cce4f5301ba604946d5e4bc53 Mon Sep 17 00:00:00 2001 From: Fabio Barone Date: Fri, 26 Apr 2019 15:34:51 -0500 Subject: [PATCH] swarm/network/stream: added pure retrieval test (syncing disabled) --- .../network/stream/snapshot_retrieval_test.go | 219 +++++++++++++++--- 1 file changed, 191 insertions(+), 28 deletions(-) diff --git a/swarm/network/stream/snapshot_retrieval_test.go b/swarm/network/stream/snapshot_retrieval_test.go index e34f87951b..8dafb0be5b 100644 --- a/swarm/network/stream/snapshot_retrieval_test.go +++ b/swarm/network/stream/snapshot_retrieval_test.go @@ -16,6 +16,7 @@ package stream import ( + "bytes" "context" "fmt" "sync" @@ -33,17 +34,17 @@ import ( "github.com/ethereum/go-ethereum/swarm/testutil" ) -//constants for random file generation +// constants for random file generation const ( minFileSize = 2 maxFileSize = 40 ) -//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. +// 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) { var nodeCount []int @@ -67,6 +68,44 @@ func TestFileRetrieval(t *testing.T) { } } +// TestPureRetrieval tests pure retrieval without syncing +// A configurable number of nodes and chunks +// can be provided to the test. +// A number of random chunks is generated, then stored directly in +// each node's localstore according to their address. +// Each chunk is supposed to end up at certain nodes +// With retrieval we then make sure that every node can actually retrieve +// the chunks. +func TestPureRetrieval(t *testing.T) { + var nodeCount []int + var chunkCount []int + + if *nodes != 0 && *chunks != 0 { + nodeCount = []int{*nodes} + chunkCount = []int{*chunks} + } else { + nodeCount = []int{16} + chunkCount = []int{16} + + if *longrunning { + nodeCount = append(nodeCount, 32, 64) + chunkCount = append(chunkCount, 32, 256) + } else if testutil.RaceEnabled { + nodeCount = []int{4} + chunkCount = []int{4} + } + + } + + for _, nc := range nodeCount { + for _, c := range chunkCount { + if err := runPureRetrievalTest(nc, c); err != nil { + t.Error(err) + } + } + } +} + //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 @@ -132,14 +171,148 @@ var retrievalSimServiceMap = map[string]simulation.ServiceFunc{ }, } -/* -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. +// runPureRetrievalTest by uploading a snapshot, +// then starting a simulation, distribute chunks to nodes +// and start retrieval. +// The snapshot should have 'streamer' in its service list. +func runPureRetrievalTest(nodeCount int, chunkCount int) error { + // the pure retrieval test needs a different service map, as we want + // syncing disabled and we don't need to set the syncUpdateDelay + sim := simulation.New(map[string]simulation.ServiceFunc{ + "streamer": func(ctx *adapters.ServiceContext, bucket *sync.Map) (s node.Service, cleanup func(), err error) { + addr, netStore, delivery, clean, err := newNetStoreAndDelivery(ctx, bucket) + if err != nil { + return nil, nil, err + } -The snapshot should have 'streamer' in its service list. -*/ + r := NewRegistry(addr.ID(), delivery, netStore, state.NewInmemoryStore(), &RegistryOptions{ + Syncing: SyncingDisabled, // disable syncing + }, nil) + + cleanup = func() { + r.Close() + clean() + } + + return r, cleanup, nil + }, + }, + ) + defer sim.Close() + + log.Info("Initializing test config", "node count", nodeCount) + + 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) + + ctx, cancelSimRun := context.WithTimeout(context.Background(), 3*time.Minute) + defer cancelSimRun() + + filename := fmt.Sprintf("testing/snapshot_%d.json", nodeCount) + err := sim.UploadSnapshot(ctx, filename) + if err != nil { + return err + } + + log.Info("Starting simulation") + + result := sim.Run(ctx, func(ctx context.Context, sim *simulation.Simulation) error { + nodeIDs := sim.UpNodeIDs() + // first iteration: create addresses + 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 + } + + // now create random chunks + chunks := make([]chunk.Chunk, chunkCount) + + for i := 0; i < chunkCount; i++ { + chunks[i] = storage.GenerateRandomChunk(int64(chunkSize)) + conf.hashes = append(conf.hashes, chunks[i].Address()) + } + log.Debug("random chunks generated, mapping keys to nodes") + + // map addresses to nodes + mapKeysToNodes(conf) + + // second iteration: store chunks at the nodes they would be + // expected to be + log.Debug("storing every chunk at correspondent node store") + for _, id := range nodeIDs { + // these are the chunks for this node + localChunks := conf.idToChunksMap[id] + // for every such chunk (which are only indexes)... + for _, ch := range localChunks { + item, ok := sim.NodeItem(id, bucketKeyStore) + if !ok { + return fmt.Errorf("Error accessing localstore") + } + lstore := item.(chunk.Store) + // ...get the actual chunk + for _, chnk := range chunks { + if bytes.Equal(chnk.Address(), conf.hashes[ch]) { + // ...and store it in the localstore + _, err = lstore.Put(ctx, chunk.ModePutUpload, chnk) + } + } + } + } + + // now try to retrieve every chunk from every node + log.Debug("starting retrieval") + cnt := 0 + + REPEAT: + for { + for _, id := range nodeIDs { + item, ok := sim.NodeItem(id, bucketKeyFileStore) + if !ok { + return fmt.Errorf("No filestore") + } + fileStore := item.(*storage.FileStore) + for _, chunk := range chunks { + reader, _ := fileStore.Retrieve(context.TODO(), chunk.Address()) + //check that we can read the file size and that it corresponds to the generated file size + if s, err := reader.Size(ctx, nil); err != nil || s != int64(chunkSize) { + log.Debug("Retrieve error", "err", err, "hash", chunk.Address(), "nodeId", id) + time.Sleep(500 * time.Millisecond) + // we continue trying if it failed + continue REPEAT + } + log.Debug(fmt.Sprintf("chunk with root hash %x successfully retrieved", chunk.Address())) + cnt++ + } + } + // we iterate sequentially, one after the other, so at this point we got all chunks + break + } + log.Info("retrieval terminated, chunks retrieved: ", "count", cnt) + return nil + + }) + + log.Info("Simulation terminated") + + if result.Error != nil { + return result.Error + } + return nil +} + +// The test loads a snapshot file to construct the swarm network. +// The snapshot should have 'streamer' in its service list. func runFileRetrievalTest(nodeCount int) error { sim := simulation.New(retrievalSimServiceMap) defer sim.Close() @@ -180,9 +353,6 @@ func runFileRetrievalTest(nodeCount int) error { //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 conf.hashes, randomFiles, err = uploadFilesToNodes(sim) if err != nil { @@ -227,16 +397,9 @@ func runFileRetrievalTest(nodeCount int) error { return nil } -/* -The test generates the given number of chunks. - -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. -*/ +// The test generates the given number of chunks. +// The test loads a snapshot file to construct the swarm network. +// The snapshot should have 'streamer' in its service list. func runRetrievalTest(t *testing.T, chunkCount int, nodeCount int) error { t.Helper() sim := simulation.New(retrievalSimServiceMap) @@ -278,8 +441,8 @@ func runRetrievalTest(t *testing.T, chunkCount int, nodeCount int) error { if !ok { return fmt.Errorf("No localstore") } - store := item.(chunk.Store) - conf.hashes, err = uploadFileToSingleNodeStore(node.ID(), chunkCount, store) + lstore := item.(chunk.Store) + conf.hashes, err = uploadFileToSingleNodeStore(node.ID(), chunkCount, lstore) if err != nil { return err }