diff --git a/swarm/network/stream/common_test.go b/swarm/network/stream/common_test.go index ef6a12bfbb..dc8a217821 100644 --- a/swarm/network/stream/common_test.go +++ b/swarm/network/stream/common_test.go @@ -50,14 +50,18 @@ 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 + waitPeerErrC chan error + chunkSize = 4096 + registries map[discover.NodeID]*TestRegistry + createStoreFunc func(id discover.NodeID, addr *network.BzzAddr) (storage.ChunkStore, error) + getRetrieveFunc = defaultRetrieveFunc + subscriptionCount = 0 ) var services = adapters.Services{ @@ -90,19 +94,24 @@ func NewStreamerService(ctx *adapters.ServiceContext) (node.Service, error) { delivery := NewDelivery(kad, db) deliveries[id] = delivery r := NewRegistry(addr, delivery, db, state.NewInmemoryStore(), &RegistryOptions{ - SkipCheck: defaultSkipCheck, + SkipCheck: defaultSkipCheck, + DoRetrieve: false, }) 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..2f9ae037d1 --- /dev/null +++ b/swarm/network/stream/snapshot_retrieval_test.go @@ -0,0 +1,794 @@ +// 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" + "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" + "github.com/ethereum/go-ethereum/swarm/network" + streamTesting "github.com/ethereum/go-ethereum/swarm/network/stream/testing" + "github.com/ethereum/go-ethereum/swarm/storage" +) + +//constants for random file generation +const ( + minFileSize = 2 + maxFileSize = 40 +) + +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) + //data directories for each node and store + 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 + return deliveries[id].RequestFromPeers(chunk.Key[:], skipCheck) + } + } + //registries, map of discover.NodeID to its streamer + registries = make(map[discover.NodeID]*TestRegistry) + //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 + } +} + +//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} + } 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 +//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} + } + for _, n := range nodeCnt { + for _, c := range chnkCnt { + retrievalTest(t, c, n) + } + } + } +} + +//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 + 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 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 + } + var rpcSubscriptionsWg sync.WaitGroup + //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.NewPeerPotMap(testMinProxBinSize, 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 + addr := common.Bytes2Hex(r.addr.OAddr) + pp := ppmap[addr] + //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 + } + 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 + wdDoneC, err := streamTesting.WatchDisconnections(id, client, disconnectC, quitC) + if err != nil { + return err + } + rpcSubscriptionsWg.Add(1) + go func() { + <-wdDoneC + rpcSubscriptionsWg.Done() + }() + } + + //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. + +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. 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 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 + } + var rpcSubscriptionsWg sync.WaitGroup + //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() + //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.addrToIdMap[string(a)] = ids[c] + } + + //needed for healthy call + ppmap = network.NewPeerPotMap(testMinProxBinSize, conf.addrs) + + 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 + addr := common.Bytes2Hex(network.ToOverlayAddr(id.Bytes())) + pp := ppmap[addr] + //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, 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(), 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 + } + + //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 + wdDoneC, err := streamTesting.WatchDisconnections(id, client, disconnectC, quitC) + if err != nil { + return err + } + rpcSubscriptionsWg.Add(1) + go func() { + <-wdDoneC + rpcSubscriptionsWg.Done() + }() + } + + //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 { + //now upload the chunks to the selected random single node + chnks, err := uploadFileToSingleNodeStore(uploadNode.ID(), chunkCount) + if err != nil { + return err + } + conf.hashes = append(conf.hashes, 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) { + + //don't check the uploader node + 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 + + //check on the node's dpa (netstore) + dpa := registries[id].dpa + //check all 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 { + 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 := 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 +} + +//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) + + 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], 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") + wait() + if err != nil { + return nil, nil, err + } + rootkeys[i] = rk + } + return rootkeys, rfiles, nil +} + +//generate a random file (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) + _, err := crand.Read(b) + if err != nil { + log.Error("Error generating random file.", "err", err) + return "", err + } + return string(b), nil +} diff --git a/swarm/network/stream/snapshot_sync_test.go b/swarm/network/stream/snapshot_sync_test.go index 08f3d500c6..06b20ff06b 100644 --- a/swarm/network/stream/snapshot_sync_test.go +++ b/swarm/network/stream/snapshot_sync_test.go @@ -43,7 +43,7 @@ import ( ) const testMinProxBinSize = 2 -const MAX_TIMEOUT = 600 +const MaxTimeout = 600 var ( pof = pot.DefaultPof(256) @@ -54,8 +54,6 @@ var ( datadirs map[discover.NodeID]string ppmap map[string]*network.PeerPot - globalWg sync.WaitGroup - live bool history bool @@ -64,10 +62,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 } @@ -85,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) @@ -95,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 @@ -107,20 +103,34 @@ 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 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} + } 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) } @@ -128,11 +138,9 @@ func TestLongRunningSyncing(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 @@ -159,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 @@ -179,20 +194,22 @@ 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 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) - //First load the snapshot from the file + //array where the generated chunk hashes will be stored + conf.hashes = make([]storage.Key, 0) + //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 @@ -235,7 +252,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") @@ -245,6 +261,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 { @@ -276,7 +294,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) } @@ -285,13 +303,21 @@ 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") - // 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: %d", j)) + log.Trace(fmt.Sprintf("Subscribe to subscription events: %d", j)) client, err := net.GetNode(id).Client() if err != nil { return err @@ -308,19 +334,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 { @@ -331,33 +344,50 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { <-wdDoneC rpcSubscriptionsWg.Done() }() - //start syncing! - err = client.CallContext(ctx, nil, "stream_startSyncing") + } + + //second iteration: start syncing + 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 + } + //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 - go func() { - globalWg.Wait() - errc <- nil - }() - - err := <-errc - if err != nil { - return err - } - log.Info("Stream subscriptions successfully requested") - if live { - //now upload the chunks to the selected random single node - chunks, err := uploadFileToSingleNodeStore(node.ID(), chunkCount) + //`watchSubscriptionEvents` will write with a `nil` value to errc + for err := range errc { if err != nil { return err } - conf.chunks = append(conf.chunks, chunks...) + //`nil` received, decrement count + subscriptionCount-- + //all subscriptions received + if subscriptionCount == 0 { + break + } + } + + log.Info("Stream subscriptions successfully requested") + if live { + //now upload the chunks to the selected random single node + hashes, err := uploadFileToSingleNodeStore(node.ID(), chunkCount) + if err != nil { + return err + } + 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)) } @@ -388,7 +418,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 { @@ -416,7 +446,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() @@ -437,13 +467,11 @@ 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 -func (r *TestRegistry) StartSyncing(ctx context.Context) error { +//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 if log.Lvl(*loglevel) == log.LvlDebug { @@ -459,9 +487,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 +500,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 +509,7 @@ func (r *TestRegistry) StartSyncing(ctx context.Context) error { return true }) - return nil + return subCnt, nil } //map chunk keys to addresses which are responsible @@ -495,11 +524,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 @@ -516,7 +545,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 @@ -616,6 +645,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 +666,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 {