diff --git a/swarm/network/stream/snapshot_sync_test.go b/swarm/network/stream/snapshot_sync_test.go index 52517e9541..e028b495eb 100644 --- a/swarm/network/stream/snapshot_sync_test.go +++ b/swarm/network/stream/snapshot_sync_test.go @@ -52,6 +52,7 @@ var ( requestedSubscriptions int receivedSubscriptions int + subscriptionsFinished bool printed bool ) @@ -94,10 +95,9 @@ func initSyncTest() { 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) - // peerCount function gives the number of peer connections for a nodeID - // this is needed for the service run function to wait until - // each protocol instance runs and the streamer peers are available + //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 @@ -210,10 +210,12 @@ func runSyncTest(chunkCount int, nodeCount int) error { disconnectC := make(chan error) quitC := make(chan struct{}) + //load nodes from the snapshot file net, err := initNetWithSnapshot(nodeCount) if err != nil { return err } + //define the cleanup function cleanup := func() { timingTicker.Stop() actionTicker.Stop() @@ -226,6 +228,7 @@ func runSyncTest(chunkCount int, nodeCount int) error { //finally clear all data directories datadirsCleanup() } + //do cleanup after test is terminated defer cleanup() //get the nodes of the network nodes := net.GetNodes() @@ -252,8 +255,10 @@ func runSyncTest(chunkCount int, nodeCount int) error { } log.Info("Test config successfully initialized") + //only needed for healthy call when debugging ppmap = network.NewPeerPot(testMinProxBinSize, ids, conf.addrs) + //variables needed to wait for all subscriptions established before uploading subscriptionsDone := make(chan struct{}) errc := make(chan error) //define the action to be performed before the test checks: start syncing @@ -262,18 +267,17 @@ func runSyncTest(chunkCount int, nodeCount int) error { // each node Subscribes to each other's swarmChunkServerStreamName for j, id := range ids { log.Trace(fmt.Sprintf("subscribe: %d", j)) - ctx, cancel := context.WithTimeout(ctx, 1*time.Second) + ctx, cancel := context.WithTimeout(context.Background(), 1*time.Second) defer cancel() client, err := net.GetNode(id).Client() if err != nil { return err } - /* - watchCtx, watchCancel := context.WithTimeout(ctx, 15*time.Second) - defer watchCancel() - watchSubscriptionEvents(watchCtx, id, client, subscriptionsDone, errc) - */ + //now setup and start event watching in order to know when we can upload + watchCtx, watchCancel := context.WithTimeout(context.Background(), 15*time.Second) + defer watchCancel() + watchSubscriptionEvents(watchCtx, id, client, subscriptionsDone, errc) if log.Lvl(*loglevel) == log.LvlDebug { //print uploading node kademlia @@ -288,16 +292,21 @@ func runSyncTest(chunkCount int, nodeCount int) error { log.Debug(kt) } } + //watch for peers disconnecting err = streamTesting.WatchDisconnections(id, client, disconnectC, quitC) if err != nil { return err } + //start syncing! err = client.CallContext(ctx, nil, "stream_startSyncing") if err != nil { return err } } - + //only at this point all subscriptions have been finished and + //and we have a final number of subscriptions we need to wait for + subscriptionsFinished = true + //now wait until the number of expected subscriptions has been finished select { case <-subscriptionsDone: close(subscriptionsDone) @@ -305,9 +314,6 @@ func runSyncTest(chunkCount int, nodeCount int) error { return err } log.Info("Stream subscriptions successfully requested") - // wait for subscritpions - //TODO: Implement a proper sync mechanism so that we don't need to Sleep() - //time.Sleep(10 * time.Second) //now upload the chunks to the selected random single node conf.chunks, err = uploadFileToSingleNodeStore(node.ID(), chunkCount) if err != nil { @@ -397,10 +403,12 @@ func runSyncTest(chunkCount int, nodeCount int) 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 { var err error @@ -432,7 +440,6 @@ func (r *TestRegistry) StartSyncing(ctx context.Context) error { return false } requestedSubscriptions += 1 - //fmt.Println(requestedSubscriptions) return true }) @@ -515,7 +522,6 @@ func mapKeysToNodes(conf *synctestConfig) { return true }) kmap[string(conf.chunks[i])] = nns - //log.Debug(fmt.Sprintf("Length for id %s: %d",ids[i],len(kmap[ids[i]]))) } if log.Lvl(*loglevel) == log.LvlTrace { for k, v := range nodemap { @@ -639,7 +645,16 @@ func initNetWithSnapshot(nodeCount int) (*simulations.Network, error) { return nil, err } + //the snapshot probably has the property EnableMsgEvents not set + //just in case, set it to true! + //(we need this to wait for messages before uploading) + for _, n := range snap.Nodes { + n.Node.Config.EnableMsgEvents = true + } + log.Info("Waiting for p2p connections to be established...") + //wait until all node connections are established + //setup variables errc := make(chan error) ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) defer cancel() @@ -647,10 +662,12 @@ func initNetWithSnapshot(nodeCount int) (*simulations.Network, error) { done := make(chan struct{}) go waitForSnapshotConnsUp(ctx, net, done, connCount, errc) + //now we can load the snapshot err = net.Load(&snap) if err != nil { return nil, err } + //finally wait until connections are established select { case <-done: close(done) @@ -661,26 +678,29 @@ func initNetWithSnapshot(nodeCount int) (*simulations.Network, error) { return net, nil } +//we want to wait for subscriptions to be established before uploading to test +//that live syncing is working correctly func watchSubscriptionEvents(ctx context.Context, id discover.NodeID, client *rpc.Client, done chan struct{}, errc chan error) { events := make(chan *p2p.PeerEvent) sub, err := client.Subscribe(context.Background(), "admin", events, "peerEvents") if err != nil { + log.Error(err.Error()) errc <- fmt.Errorf("error getting peer events for node %v: %s", id, err) return } go func() { + defer sub.Unsubscribe() + for { select { case <-ctx.Done(): return case e := <-events: - fmt.Println(e) - fmt.Println(*e.MsgCode) - if e.Type == p2p.PeerEventTypeMsgRecv && e.Protocol == "stream" && e.MsgCode != nil && *e.MsgCode == 1 { - fmt.Println(receivedSubscriptions) - fmt.Println(requestedSubscriptions) + //just catch SubscribeMsg + if e.Type == p2p.PeerEventTypeMsgRecv && e.Protocol == "stream" && e.MsgCode != nil && *e.MsgCode == 4 { receivedSubscriptions += 1 - if receivedSubscriptions == requestedSubscriptions { + //only check for done if subscription process is finished + if subscriptionsFinished && (receivedSubscriptions == requestedSubscriptions) { done <- struct{}{} return }