diff --git a/swarm/network/stream/common_test.go b/swarm/network/stream/common_test.go index 8a8b7bce6a..241bd7fc3b 100644 --- a/swarm/network/stream/common_test.go +++ b/swarm/network/stream/common_test.go @@ -48,8 +48,8 @@ var ( stores map[discover.NodeID]storage.ChunkStore toAddr func(discover.NodeID) *network.BzzAddr peerCount func(discover.NodeID) int - adapter = flag.String("adapter", "socket", "type of simulation: sim|socket|exec|docker") - loglevel = flag.Int("loglevel", 4, "verbosity of logs") + adapter = flag.String("adapter", "sim", "type of simulation: sim|socket|exec|docker") + loglevel = flag.Int("loglevel", 2, "verbosity of logs") ) var ( diff --git a/swarm/network/stream/delivery.go b/swarm/network/stream/delivery.go index 33449780a4..9f0109fefd 100644 --- a/swarm/network/stream/delivery.go +++ b/swarm/network/stream/delivery.go @@ -127,7 +127,7 @@ type RetrieveRequestMsg struct { } func (d *Delivery) handleRetrieveRequestMsg(sp *Peer, req *RetrieveRequestMsg) error { - log.Debug("received request", "peer", sp.ID(), "hash", req.Key) + log.Trace("received request", "peer", sp.ID(), "hash", req.Key) s, err := sp.getServer(NewStream(swarmChunkServerStreamName, "", false)) if err != nil { return err @@ -157,6 +157,7 @@ func (d *Delivery) handleRetrieveRequestMsg(sp *Peer, req *RetrieveRequestMsg) e if req.SkipCheck { err := sp.Deliver(chunk, s.priority) if err != nil { + log.Warn("ERROR in handleRetrieveRequestMsg, DROPPING peer!", "err", err) sp.Drop(err) } } diff --git a/swarm/network/stream/messages.go b/swarm/network/stream/messages.go index 3d18a321e7..3e5c03d4d0 100644 --- a/swarm/network/stream/messages.go +++ b/swarm/network/stream/messages.go @@ -110,7 +110,7 @@ func (p *Peer) handleSubscribeMsg(req *SubscribeMsg) (err error) { go func() { if err := p.SendOfferedHashes(os, from, to); err != nil { - log.Error("ERROR in SendOfferedHashes, DROPPING peer!", "err", err) + log.Warn("ERROR in SendOfferedHashes, DROPPING peer!", "err", err) p.Drop(err) } }() @@ -128,7 +128,7 @@ func (p *Peer) handleSubscribeMsg(req *SubscribeMsg) (err error) { } go func() { if err := p.SendOfferedHashes(os, req.History.From, req.History.To); err != nil { - log.Error("ERROR in SendOfferedHashes, DROPPING peer!", "err", err) + log.Warn("ERROR in SendOfferedHashes, DROPPING peer!", "err", err) p.Drop(err) } }() @@ -238,12 +238,12 @@ func (p *Peer) handleOfferedHashesMsg(req *OfferedHashesMsg) error { go func() { select { case <-time.After(120 * time.Second): - log.Error("ERROR in handleOfferedHashesMsg, DROPPING peer!", "err", "TIMEOUT") + log.Warn("ERROR in handleOfferedHashesMsg, DROPPING peer!", "err", "TIMEOUT") p.Drop(err) return case err := <-c.next: if err != nil { - log.Error("ERROR in handleOfferedHashesMsg, DROPPING peer!", "err", err) + log.Warn("ERROR in handleOfferedHashesMsg, DROPPING peer!", "err", err) p.Drop(err) return } @@ -251,7 +251,7 @@ func (p *Peer) handleOfferedHashesMsg(req *OfferedHashesMsg) error { log.Trace("sending want batch", "peer", p.ID(), "stream", msg.Stream, "from", msg.From, "to", msg.To) err := p.SendPriority(msg, c.priority) if err != nil { - log.Error("ERROR in handleOfferedHashesMsg, DROPPING peer!", "err", err) + log.Warn("ERROR in handleOfferedHashesMsg, DROPPING peer!", "err", err) p.Drop(err) } }() @@ -284,7 +284,7 @@ func (p *Peer) handleWantedHashesMsg(req *WantedHashesMsg) error { // launch in go routine since GetBatch blocks until new hashes arrive go func() { if err := p.SendOfferedHashes(s, req.From, req.To); err != nil { - log.Error("ERROR in handleWantedHashesMsg, DROPPING peer!", "err", err) + log.Warn("ERROR in handleWantedHashesMsg, DROPPING peer!", "err", err) p.Drop(err) } }() diff --git a/swarm/network/stream/snapshot_sync_test.go b/swarm/network/stream/snapshot_sync_test.go index b4c8744ea7..fb54641c67 100644 --- a/swarm/network/stream/snapshot_sync_test.go +++ b/swarm/network/stream/snapshot_sync_test.go @@ -19,6 +19,7 @@ import ( "context" crand "crypto/rand" "encoding/json" + "flag" "fmt" "io" "io/ioutil" @@ -52,14 +53,17 @@ var ( datadirs map[discover.NodeID]string ppmap map[discover.NodeID]*network.PeerPot + globalWg sync.WaitGroup + live bool history bool + + longrunning = flag.Bool("longrunning", false, "do run long-running tests") ) type synctestConfig struct { addrs [][]byte chunks []storage.Key - retrievalMap map[string]map[string]time.Duration idToChunksMap map[discover.NodeID][]int chunksToNodesMap map[string][]int idToAddrMap map[discover.NodeID][]byte @@ -68,8 +72,6 @@ type synctestConfig struct { func init() { rand.Seed(time.Now().Unix()) - - //initSyncTest() } //common_test needs to initialize the test in a init() func @@ -80,7 +82,6 @@ func initSyncTest() { //assign the toAddr func so NewStreamerService can build the addr toAddr = func(id discover.NodeID) *network.BzzAddr { addr := network.NewAddrFromNodeID(id) - // addr.OAddr[0] = byte(0) return addr } @@ -93,7 +94,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) @@ -110,46 +110,24 @@ func initSyncTest() { //TestSyncing_x_y //x is the number of chunks which will be uploaded //y is the number of nodes for the test -func TestSyncing_1_16(t *testing.T) { testSyncing(t, 1, 16) } -func TestSyncing_1_32(t *testing.T) { testSyncing(t, 1, 32) } -func TestSyncing_1_64(t *testing.T) { testSyncing(t, 1, 64) } -func TestSyncing_1_128(t *testing.T) { testSyncing(t, 1, 128) } -func TestSyncing_1_256(t *testing.T) { testSyncing(t, 1, 256) } -func TestSyncing_4_16(t *testing.T) { testSyncing(t, 4, 16) } -func TestSyncing_4_32(t *testing.T) { testSyncing(t, 4, 32) } -func TestSyncing_4_64(t *testing.T) { testSyncing(t, 4, 64) } -func TestSyncing_4_128(t *testing.T) { testSyncing(t, 4, 128) } -func TestSyncing_8_16(t *testing.T) { testSyncing(t, 8, 16) } -func TestSyncing_8_32(t *testing.T) { testSyncing(t, 8, 32) } -func TestSyncing_8_64(t *testing.T) { testSyncing(t, 8, 64) } -func TestSyncing_8_128(t *testing.T) { testSyncing(t, 8, 128) } -func TestSyncing_32_16(t *testing.T) { testSyncing(t, 32, 16) } -func TestSyncing_32_32(t *testing.T) { testSyncing(t, 32, 32) } -func TestSyncing_32_64(t *testing.T) { testSyncing(t, 32, 64) } -func TestSyncing_128_16(t *testing.T) { testSyncing(t, 128, 16) } -func TestSyncing_128_32(t *testing.T) { testSyncing(t, 128, 32) } -func TestSyncing_128_64(t *testing.T) { testSyncing(t, 128, 64) } -func TestSyncing_256_16(t *testing.T) { testSyncing(t, 256, 16) } -func TestSyncing_256_32(t *testing.T) { testSyncing(t, 256, 32) } -func TestSyncing_1024_16(t *testing.T) { testSyncing(t, 1024, 16) } +func TestSyncing_4_32(t *testing.T) { testSyncing(t, 4, 32) } +func TestSyncing_32_16(t *testing.T) { testSyncing(t, 32, 16) } -//The following tests have been disabled because they seem to hit resource limits -//on developer machines and/or are long running -/* -func TestSyncing_256_64(t *testing.T) { testSyncing(t, 256, 64) } -func TestSyncing_8_256(t *testing.T) { testSyncing(t, 8, 256) } -func TestSyncing_32_256(t *testing.T) { testSyncing(t, 32, 256) } -func TestSyncing_32_128(t *testing.T) { testSyncing(t, 32, 128) } -func TestSyncing_256_128(t *testing.T) { testSyncing(t, 256, 128) } -func TestSyncing_128_128(t *testing.T) { testSyncing(t, 128, 128) } -func TestSyncing_256_256(t *testing.T) { testSyncing(t, 256, 256) } -func TestSyncing_128_256(t *testing.T) { testSyncing(t, 128, 256) } -func TestSyncing_1024_32(t *testing.T) { testSyncing(t, 1024, 32) } -func TestSyncing_1024_64(t *testing.T) { testSyncing(t, 1024, 64) } -func TestSyncing_1024_128(t *testing.T) { testSyncing(t, 1024, 128) } -func TestSyncing_1024_256(t *testing.T) { testSyncing(t, 1024, 256) } -*/ +func TestLongRunningSyncing(t *testing.T) { + if *longrunning { + chnkCnt := []int{1, 8, 32, 256, 1024} + nCnt := []int{16, 32, 64, 128, 256} + for _, chnk := range chnkCnt { + for _, n := range nCnt { + log.Info(fmt.Sprintf("Long running test with %d chunks and %d nodes...", chnk, n)) + testSyncing(t, chnk, n) + } + } + } +} + +//do run the tests func testSyncing(t *testing.T, chunkCount int, nodeCount int) { initSyncTest() ids = make([]discover.NodeID, nodeCount) @@ -162,7 +140,7 @@ func testSyncing(t *testing.T, chunkCount int, nodeCount int) { if err != nil { t.Fatal(err) } - //finaly test history only + //test history only log.Info("Testing history only") live = false history = true @@ -170,7 +148,7 @@ func testSyncing(t *testing.T, chunkCount int, nodeCount int) { if err != nil { t.Fatal(err) } - //test live and history + //finally test live and history log.Info("Testing live and history") live = true err = runSyncTest(chunkCount, nodeCount, live, history) @@ -202,8 +180,6 @@ For every test run, a series of three tests will be executed: func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { //initialize the test struct conf = &synctestConfig{} - //mapping of nearest node addresses for chunk hashes - conf.retrievalMap = make(map[string]map[string]time.Duration) //map of discover ID to indexes of chunks expected at that ID conf.idToChunksMap = make(map[discover.NodeID][]int) //map of discover ID to kademlia overlay address @@ -231,7 +207,6 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { localStoreCleanup() //finally clear all data directories datadirsCleanup() - //close(disconnectC) close(quitC) }() //get the nodes of the network @@ -275,8 +250,8 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { //call Healthy RPC h := r.delivery.overlay.Healthy(pp) //print info - log.Error(r.delivery.overlay.String()) - log.Error(fmt.Sprintf("IS HEALTHY: %t", h.GotNN && h.KnowNN && h.Full)) + 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 @@ -299,9 +274,8 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { mapKeysToNodes(conf) } - errc := make(chan error) //variables needed to wait for all subscriptions established before uploading - subscriptionsDone := make(chan struct{}) + 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) @@ -309,20 +283,19 @@ 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 - var wg sync.WaitGroup for j, id := range ids { log.Trace(fmt.Sprintf("subscribe: %d", j)) client, err := net.GetNode(id).Client() if err != nil { return err } - wg.Add(1) - watchSubscriptionEvents(ctx, id, client, &wg, errc) - if log.Lvl(*loglevel) >= log.LvlDebug { - //uncomment this if to see only the uploader's node kademlia - //otherwise print all kademlias - //if j == idx { + watchSubscriptionEvents(ctx, id, client, errc) + + 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 { @@ -331,7 +304,6 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { log.Debug("kad table " + node.ID().String()) log.Debug(kt) - //} } //watch for peers disconnecting err = streamTesting.WatchDisconnections(id, client, disconnectC, quitC) @@ -346,16 +318,16 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { } //now wait until the number of expected subscriptions has been finished - go func() { - wg.Wait() - close(subscriptionsDone) + globalWg.Wait() + errc <- nil }() select { - case <-subscriptionsDone: case err := <-errc: - return err + if err != nil { + return err + } } log.Info("Stream subscriptions successfully requested") if live { @@ -411,10 +383,6 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { return allSuccess, nil } - timeout := MAX_TIMEOUT * time.Second - ctx, cancel := context.WithTimeout(context.Background(), timeout) - defer cancel() - //for each tick, run the checks on all nodes timingTicker := time.NewTicker(time.Second * 1) defer timingTicker.Stop() @@ -428,6 +396,11 @@ func runSyncTest(chunkCount int, nodeCount int, live bool, history bool) error { }() 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, @@ -475,13 +448,12 @@ func (r *TestRegistry) StartSyncing(ctx context.Context) error { 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 log.Debug(fmt.Sprintf("Requesting subscription by: registry %s from peer %s for bin: %d", r.addr.ID(), conf.addrToIdMap[string(conn.Address())], po)) - //fmt.Printf("rs: %s peer %s bin %d\n", r.addr.ID(), conf.addrToIdMap[string(conn.Address())], po) var histRange *Range - if !history { - histRange = nil - } else { + if history { histRange = &Range{} } + + globalWg.Add(1) 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)) @@ -528,21 +500,6 @@ func mapKeysToNodes(conf *synctestConfig) { }) kmap[string(conf.chunks[i])] = nns } - if log.Lvl(*loglevel) == log.LvlTrace { - for k, v := range nodemap { - log.Trace(fmt.Sprintf("Node %s: ", conf.addrToIdMap[k])) - for _, vv := range v { - log.Trace(fmt.Sprintf("%v", conf.chunks[vv])) - } - log.Trace(fmt.Sprintf("%v", conf.addrToIdMap[k])) - } - for k, v := range kmap { - log.Trace(fmt.Sprintf("Chunk %s: ", k)) - for _, vv := range v { - log.Trace(fmt.Sprintf("%v", conf.addrToIdMap[string(conf.addrs[vv])])) - } - } - } for addr, chunks := range nodemap { //this selects which chunks are expected to be found with the given node conf.idToChunksMap[conf.addrToIdMap[addr]] = chunks @@ -556,11 +513,10 @@ func uploadFileToSingleNodeStore(id discover.NodeID, chunkCount int) ([]storage. log.Debug(fmt.Sprintf("Uploading to node id: %s", id)) lstore := stores[id] size := chunkSize - dpa := storage.NewDPA(lstore, storage.NewChunkerParams()) - dpa.Start() + dpa := storage.NewDPA(lstore, storage.NewDPAParams()) var rootkeys []storage.Key for i := 0; i < chunkCount; i++ { - rk, wait, err := dpa.Store(io.LimitReader(crand.Reader, int64(size)), int64(size)) + rk, wait, err := dpa.Store(io.LimitReader(crand.Reader, int64(size)), int64(size), false) wait() if err != nil { return nil, err @@ -568,42 +524,9 @@ func uploadFileToSingleNodeStore(id discover.NodeID, chunkCount int) ([]storage. rootkeys = append(rootkeys, (rk)) } - defer dpa.Stop() - return rootkeys, nil } -//Here we wait until all connections from the snapshot are up -func waitForSnapshotConnsUp(ctx context.Context, net *simulations.Network, connCount int, errc chan error) { - arrivedConns := 0 - events := make(chan *simulations.Event) - //subscribe to all events from the network - sub := net.Events().Subscribe(events) - defer sub.Unsubscribe() - - for { - select { - case <-ctx.Done(): - errc <- fmt.Errorf("Timeout waiting for Snapshot connections") - case event := <-events: - //if the event is of type connection, is a Live event and the connection is up - //NOTE; this will require that all connections are UP in the snapshot! - if event.Type == simulations.EventTypeConn && !event.Control && event.Conn.Up { - arrivedConns++ - //the amount of expected connections has been reached, so we can stop waiting - if arrivedConns == connCount { - errc <- nil - return - } - } - case err := <-sub.Err(): - if err != nil { - errc <- err - } - } - } -} - //initialize a network from a snapshot func initNetWithSnapshot(nodeCount int) (*simulations.Network, error) { @@ -654,33 +577,19 @@ func initNetWithSnapshot(nodeCount int) (*simulations.Network, error) { } log.Info("Waiting for p2p connections to be established...") - //wait until all node connections are established - //setup variables - ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) - defer cancel() - connCount := len(snap.Conns) - errc := make(chan error) - go waitForSnapshotConnsUp(ctx, net, 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 err = <-errc: - if err != nil { - return nil, err - } - } - log.Info("Snapshot loaded and connections established") + log.Info("Snapshot loaded") 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, wg *sync.WaitGroup, errc chan error) { +func watchSubscriptionEvents(ctx context.Context, id discover.NodeID, client *rpc.Client, errc chan error) { events := make(chan *p2p.PeerEvent) sub, err := client.Subscribe(context.Background(), "admin", events, "peerEvents") if err != nil { @@ -694,12 +603,12 @@ func watchSubscriptionEvents(ctx context.Context, id discover.NodeID, client *rp for { select { case <-ctx.Done(): + errc <- ctx.Err() return case e := <-events: //just catch SubscribeMsg if e.Type == p2p.PeerEventTypeMsgRecv && e.Protocol == "stream" && e.MsgCode != nil && *e.MsgCode == 4 { - wg.Done() - return + globalWg.Done() } case err := <-sub.Err(): if err != nil { diff --git a/swarm/network/stream/stream.go b/swarm/network/stream/stream.go index 96f4524161..b36ad27138 100644 --- a/swarm/network/stream/stream.go +++ b/swarm/network/stream/stream.go @@ -93,7 +93,7 @@ func NewRegistry(addr *network.BzzAddr, delivery *Delivery, db *storage.DBAPI, i return NewSwarmChunkServer(delivery.db), nil }) streamer.RegisterClientFunc(swarmChunkServerStreamName, func(p *Peer, _ string, _ bool) (Client, error) { - return NewSwarmSyncerClient(p, delivery.db) + return NewSwarmSyncerClient(p, delivery.db, false) }) RegisterSwarmSyncerServer(streamer, db) RegisterSwarmSyncerClient(streamer, db) diff --git a/swarm/network/stream/syncer.go b/swarm/network/stream/syncer.go index 154fb89c04..799d097af4 100644 --- a/swarm/network/stream/syncer.go +++ b/swarm/network/stream/syncer.go @@ -135,15 +135,17 @@ type SwarmSyncerClient struct { storeC chan *storage.Chunk db *storage.DBAPI // chunker storage.Chunker - currentRoot storage.Key - requestFunc func(chunk *storage.Chunk) - end, start uint64 + currentRoot storage.Key + requestFunc func(chunk *storage.Chunk) + end, start uint64 + ignoreExistingRequest bool } // NewSwarmSyncerClient is a contructor for provable data exchange syncer -func NewSwarmSyncerClient(_ *Peer, db *storage.DBAPI) (*SwarmSyncerClient, error) { +func NewSwarmSyncerClient(_ *Peer, db *storage.DBAPI, ignoreExistingRequest bool) (*SwarmSyncerClient, error) { return &SwarmSyncerClient{ db: db, + ignoreExistingRequest: ignoreExistingRequest, }, nil } @@ -187,7 +189,7 @@ func NewSwarmSyncerClient(_ *Peer, db *storage.DBAPI) (*SwarmSyncerClient, error // to handle incoming sync streams func RegisterSwarmSyncerClient(streamer *Registry, db *storage.DBAPI) { streamer.RegisterClientFunc("SYNC", func(p *Peer, _ string, love bool) (Client, error) { - return NewSwarmSyncerClient(p, db) + return NewSwarmSyncerClient(p, db, true) }) } @@ -195,7 +197,7 @@ func RegisterSwarmSyncerClient(streamer *Registry, db *storage.DBAPI) { func (s *SwarmSyncerClient) NeedData(key []byte) (wait func()) { chunk, created := s.db.GetOrCreateRequest(key) // TODO: we may want to request from this peer anyway even if the request exists - if chunk.ReqC == nil || created == false { + if chunk.ReqC == nil || (s.ignoreExistingRequest && !created) { return nil } // create request and wait until the chunk data arrives and is stored diff --git a/swarm/network/stream/syncer_test.go b/swarm/network/stream/syncer_test.go index ad4d9634c0..c2f679e594 100644 --- a/swarm/network/stream/syncer_test.go +++ b/swarm/network/stream/syncer_test.go @@ -45,8 +45,11 @@ func TestSyncerSimulation(t *testing.T) { func testSyncBetweenNodes(t *testing.T, nodes, conns, chunkCount int, skipCheck bool, po uint8) { defaultSkipCheck = skipCheck + createStoreFunc = createTestLocalStorageFromSim + registries = make(map[discover.NodeID]*TestRegistry) toAddr = func(id discover.NodeID) *network.BzzAddr { addr := network.NewAddrFromNodeID(id) + //hack to put addresses in same space addr.OAddr[0] = byte(0) return addr }