From b151dc25a6c7b9463c8c0021558c78240f4b6c3a Mon Sep 17 00:00:00 2001 From: Janos Guljas Date: Mon, 11 Feb 2019 16:33:55 +0100 Subject: [PATCH] swarm/network/stream: terminate disconnect goruotines in tests --- swarm/network/stream/delivery_test.go | 32 ++++++++++++++----- swarm/network/stream/intervals_test.go | 13 ++++++-- swarm/network/stream/snapshot_sync_test.go | 13 ++++++-- swarm/network/stream/syncer_test.go | 13 ++++++-- .../visualized_snapshot_sync_sim_test.go | 17 +++++++--- 5 files changed, 67 insertions(+), 21 deletions(-) diff --git a/swarm/network/stream/delivery_test.go b/swarm/network/stream/delivery_test.go index d8644df6b4..13c82206b6 100644 --- a/swarm/network/stream/delivery_test.go +++ b/swarm/network/stream/delivery_test.go @@ -485,7 +485,8 @@ func testDeliveryFromNodes(t *testing.T, nodes, chunkCount int, skipCheck bool) } log.Info("Starting simulation") - ctx := context.Background() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() result := sim.Run(ctx, func(ctx context.Context, sim *simulation.Simulation) (err error) { nodeIDs := sim.UpNodeIDs() //determine the pivot node to be the first node of the simulation @@ -557,9 +558,16 @@ func testDeliveryFromNodes(t *testing.T, nodes, chunkCount int, skipCheck bool) var disconnected atomic.Value go func() { - for d := range disconnections { - if d.Error != nil { - log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID) + for { + select { + case <-ctx.Done(): + return + case d := <-disconnections: + if d.Error != nil { + log.Error("peer drop event error", "node", d.NodeID, "peer", d.PeerID, "err", err) + } else { + log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID) + } disconnected.Store(true) } } @@ -657,7 +665,8 @@ func benchmarkDeliveryFromNodes(b *testing.B, nodes, chunkCount int, skipCheck b b.Fatal(err) } - ctx := context.Background() + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() result := sim.Run(ctx, func(ctx context.Context, sim *simulation.Simulation) (err error) { nodeIDs := sim.UpNodeIDs() node := nodeIDs[len(nodeIDs)-1] @@ -687,9 +696,16 @@ func benchmarkDeliveryFromNodes(b *testing.B, nodes, chunkCount int, skipCheck b var disconnected atomic.Value go func() { - for d := range disconnections { - if d.Error != nil { - log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID) + for { + select { + case <-ctx.Done(): + return + case d := <-disconnections: + if d.Error != nil { + log.Error("peer drop event error", "node", d.NodeID, "peer", d.PeerID, "err", err) + } else { + log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID) + } disconnected.Store(true) } } diff --git a/swarm/network/stream/intervals_test.go b/swarm/network/stream/intervals_test.go index 9bd1bbd43e..cf77ebb0a1 100644 --- a/swarm/network/stream/intervals_test.go +++ b/swarm/network/stream/intervals_test.go @@ -148,9 +148,16 @@ func testIntervals(t *testing.T, live bool, history *Range, skipCheck bool) { var disconnected atomic.Value go func() { - for d := range disconnections { - if d.Error != nil { - log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID) + for { + select { + case <-ctx.Done(): + return + case d := <-disconnections: + if d.Error != nil { + log.Error("peer drop event error", "node", d.NodeID, "peer", d.PeerID, "err", err) + } else { + log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID) + } disconnected.Store(true) } } diff --git a/swarm/network/stream/snapshot_sync_test.go b/swarm/network/stream/snapshot_sync_test.go index c32ed7d077..7a1e5e51b0 100644 --- a/swarm/network/stream/snapshot_sync_test.go +++ b/swarm/network/stream/snapshot_sync_test.go @@ -164,9 +164,16 @@ func testSyncingViaGlobalSync(t *testing.T, chunkCount int, nodeCount int) { var disconnected atomic.Value go func() { - for d := range disconnections { - if d.Error != nil { - log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID) + for { + select { + case <-ctx.Done(): + return + case d := <-disconnections: + if d.Error != nil { + log.Error("peer drop event error", "node", d.NodeID, "peer", d.PeerID, "err", err) + } else { + log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID) + } disconnected.Store(true) } } diff --git a/swarm/network/stream/syncer_test.go b/swarm/network/stream/syncer_test.go index 9d96ca65ae..87df3e334d 100644 --- a/swarm/network/stream/syncer_test.go +++ b/swarm/network/stream/syncer_test.go @@ -122,9 +122,16 @@ func testSyncBetweenNodes(t *testing.T, nodes, chunkCount int, skipCheck bool, p var disconnected atomic.Value go func() { - for d := range disconnections { - if d.Error != nil { - log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID) + for { + select { + case <-ctx.Done(): + return + case d := <-disconnections: + if d.Error != nil { + log.Error("peer drop event error", "node", d.NodeID, "peer", d.PeerID, "err", err) + } else { + log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID) + } disconnected.Store(true) } } diff --git a/swarm/network/stream/visualized_snapshot_sync_sim_test.go b/swarm/network/stream/visualized_snapshot_sync_sim_test.go index c8cb5d8b38..399c0d12ad 100644 --- a/swarm/network/stream/visualized_snapshot_sync_sim_test.go +++ b/swarm/network/stream/visualized_snapshot_sync_sim_test.go @@ -81,10 +81,19 @@ func watchSim(sim *simulation.Simulation) (context.Context, context.CancelFunc) ) go func() { - for d := range disconnections { - log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID) - panic("unexpected disconnect") - cancelSimRun() + for { + select { + case <-ctx.Done(): + return + case d := <-disconnections: + if d.Error != nil { + log.Error("peer drop event error", "node", d.NodeID, "peer", d.PeerID, "err", err) + } else { + log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID) + } + panic("unexpected disconnect") + cancelSimRun() + } } }()