swarm/network/stream: add watchDisconnections helper function

This commit is contained in:
Janos Guljas 2019-02-12 11:02:59 +01:00
parent 4f5807be94
commit dee2145a51
5 changed files with 33 additions and 117 deletions

View file

@ -320,3 +320,30 @@ func createTestLocalStorageForID(id enode.ID, addr *network.BzzAddr) (storage.Ch
} }
return store, datadir, nil return store, datadir, nil
} }
// watchDisconnections receives simulation peer events in a new goroutine and sets atomic value
// disconnected to true in case of a disconnect event.
func watchDisconnections(ctx context.Context, sim *simulation.Simulation) (disconnected atomic.Value) {
log.Debug("Watching for disconnections")
disconnections := sim.PeerEvents(
context.Background(),
sim.NodeIDs(),
simulation.NewPeerEventsFilter().Drop(),
)
go func() {
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", d.Error)
} else {
log.Error("peer drop", "node", d.NodeID, "peer", d.PeerID)
}
disconnected.Store(true)
}
}
}()
return disconnected
}

View file

@ -22,7 +22,6 @@ import (
"errors" "errors"
"fmt" "fmt"
"sync" "sync"
"sync/atomic"
"testing" "testing"
"time" "time"
@ -549,29 +548,7 @@ func testDeliveryFromNodes(t *testing.T, nodes, chunkCount int, skipCheck bool)
retErrC <- err retErrC <- err
}() }()
log.Debug("Watching for disconnections") disconnected := watchDisconnections(ctx, sim)
disconnections := sim.PeerEvents(
context.Background(),
sim.NodeIDs(),
simulation.NewPeerEventsFilter().Drop(),
)
var disconnected atomic.Value
go func() {
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)
}
}
}()
defer func() { defer func() {
if err != nil { if err != nil {
if yes, ok := disconnected.Load().(bool); ok && yes { if yes, ok := disconnected.Load().(bool); ok && yes {
@ -688,28 +665,7 @@ func benchmarkDeliveryFromNodes(b *testing.B, nodes, chunkCount int, skipCheck b
return err return err
} }
disconnections := sim.PeerEvents( disconnected := watchDisconnections(ctx, sim)
context.Background(),
sim.NodeIDs(),
simulation.NewPeerEventsFilter().Drop(),
)
var disconnected atomic.Value
go func() {
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)
}
}
}()
defer func() { defer func() {
if err != nil { if err != nil {
if yes, ok := disconnected.Load().(bool); ok && yes { if yes, ok := disconnected.Load().(bool); ok && yes {

View file

@ -22,7 +22,6 @@ import (
"errors" "errors"
"fmt" "fmt"
"sync" "sync"
"sync/atomic"
"testing" "testing"
"time" "time"
@ -134,34 +133,12 @@ func testIntervals(t *testing.T, live bool, history *Range, skipCheck bool) {
liveErrC := make(chan error) liveErrC := make(chan error)
historyErrC := make(chan error) historyErrC := make(chan error)
log.Debug("Watching for disconnections")
disconnections := sim.PeerEvents(
context.Background(),
sim.NodeIDs(),
simulation.NewPeerEventsFilter().Drop(),
)
err = registry.Subscribe(storer, NewStream(externalStreamName, "", live), history, Top) err = registry.Subscribe(storer, NewStream(externalStreamName, "", live), history, Top)
if err != nil { if err != nil {
return err return err
} }
var disconnected atomic.Value disconnected := watchDisconnections(ctx, sim)
go func() {
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)
}
}
}()
defer func() { defer func() {
if err != nil { if err != nil {
if yes, ok := disconnected.Load().(bool); ok && yes { if yes, ok := disconnected.Load().(bool); ok && yes {

View file

@ -22,7 +22,6 @@ import (
"os" "os"
"runtime" "runtime"
"sync" "sync"
"sync/atomic"
"testing" "testing"
"time" "time"
@ -182,28 +181,7 @@ func testSyncingViaGlobalSync(t *testing.T, chunkCount int, nodeCount int) {
t.Fatal(err) t.Fatal(err)
} }
disconnections := sim.PeerEvents( disconnected := watchDisconnections(ctx, sim)
context.Background(),
sim.NodeIDs(),
simulation.NewPeerEventsFilter().Drop(),
)
var disconnected atomic.Value
go func() {
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)
}
}
}()
result := runSim(conf, ctx, sim, chunkCount) result := runSim(conf, ctx, sim, chunkCount)

View file

@ -24,7 +24,6 @@ import (
"math" "math"
"os" "os"
"sync" "sync"
"sync/atomic"
"testing" "testing"
"time" "time"
@ -140,28 +139,7 @@ func testSyncBetweenNodes(t *testing.T, nodes, chunkCount int, skipCheck bool, p
nodeIndex[id] = i nodeIndex[id] = i
} }
disconnections := sim.PeerEvents( disconnected := watchDisconnections(ctx, sim)
context.Background(),
sim.NodeIDs(),
simulation.NewPeerEventsFilter().Drop(),
)
var disconnected atomic.Value
go func() {
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)
}
}
}()
defer func() { defer func() {
if err != nil { if err != nil {
if yes, ok := disconnected.Load().(bool); ok && yes { if yes, ok := disconnected.Load().(bool); ok && yes {