mirror of
https://github.com/ethereum/go-ethereum.git
synced 2026-08-19 18:32:23 +00:00
swarm: PR fixes
This commit is contained in:
parent
0f8251780e
commit
577f8b72dc
1 changed files with 15 additions and 2 deletions
|
|
@ -947,6 +947,9 @@ func TestGetSubscriptionsRPC(t *testing.T) {
|
||||||
syncUpdateDelay := 500 * time.Millisecond
|
syncUpdateDelay := 500 * time.Millisecond
|
||||||
//we will later need the kad table for each node
|
//we will later need the kad table for each node
|
||||||
bucketKeyKad := simulation.BucketKey("kademlia")
|
bucketKeyKad := simulation.BucketKey("kademlia")
|
||||||
|
//holds the msg code for SubscribeMsg
|
||||||
|
var subscribeMsgCode uint64
|
||||||
|
var ok bool
|
||||||
//create a standard sim
|
//create a standard sim
|
||||||
sim := simulation.New(map[string]simulation.ServiceFunc{
|
sim := simulation.New(map[string]simulation.ServiceFunc{
|
||||||
"streamer": func(ctx *adapters.ServiceContext, bucket *sync.Map) (s node.Service, cleanup func(), err error) {
|
"streamer": func(ctx *adapters.ServiceContext, bucket *sync.Map) (s node.Service, cleanup func(), err error) {
|
||||||
|
|
@ -972,6 +975,11 @@ func TestGetSubscriptionsRPC(t *testing.T) {
|
||||||
Syncing: SyncingAutoSubscribe, //enable sync registrations
|
Syncing: SyncingAutoSubscribe, //enable sync registrations
|
||||||
SyncUpdateDelay: syncUpdateDelay,
|
SyncUpdateDelay: syncUpdateDelay,
|
||||||
}, nil)
|
}, nil)
|
||||||
|
//get the SubscribeMsg code
|
||||||
|
subscribeMsgCode, ok = r.GetSpec().GetCode(SubscribeMsg{})
|
||||||
|
if !ok {
|
||||||
|
t.Fatal("Message code for SubscribeMsg not found")
|
||||||
|
}
|
||||||
|
|
||||||
bucket.Store(bucketKeyRegistry, r)
|
bucket.Store(bucketKeyRegistry, r)
|
||||||
cleanup = func() {
|
cleanup = func() {
|
||||||
|
|
@ -1003,7 +1011,7 @@ func TestGetSubscriptionsRPC(t *testing.T) {
|
||||||
msgs := sim.PeerEvents(
|
msgs := sim.PeerEvents(
|
||||||
context.Background(),
|
context.Background(),
|
||||||
sim.NodeIDs(),
|
sim.NodeIDs(),
|
||||||
simulation.NewPeerEventsFilter().ReceivedMessages().Protocol("stream").MsgCode(4), //4 is SubscribeMsg
|
simulation.NewPeerEventsFilter().ReceivedMessages().Protocol("stream").MsgCode(subscribeMsgCode),
|
||||||
)
|
)
|
||||||
|
|
||||||
//setup the vars we need
|
//setup the vars we need
|
||||||
|
|
@ -1048,7 +1056,12 @@ func TestGetSubscriptionsRPC(t *testing.T) {
|
||||||
}
|
}
|
||||||
log.Debug("Expected message count: ", "expectedMsgCount", expectedMsgCount)
|
log.Debug("Expected message count: ", "expectedMsgCount", expectedMsgCount)
|
||||||
//wait until all subscriptions are done
|
//wait until all subscriptions are done
|
||||||
<-allSubscriptionsDone
|
select {
|
||||||
|
case <-allSubscriptionsDone:
|
||||||
|
case <-ctx.Done():
|
||||||
|
t.Fatal("Context timed out")
|
||||||
|
}
|
||||||
|
|
||||||
log.Info("All subscriptions received")
|
log.Info("All subscriptions received")
|
||||||
//now iterate again, this time we call each node via RPC to get its subscriptions
|
//now iterate again, this time we call each node via RPC to get its subscriptions
|
||||||
for _, node := range nodes {
|
for _, node := range nodes {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue