Files
Igor Sirotin fc894011a8 fix: keep allocating test ports, and move that out of the library
Removing GetFreePortIfNeeded also removed the allocation StartWakuNode
did, and the library does not treat a zero DiscV5 UDP port as "pick one":
every node tried the same default and all but the first failed to bind.

StartWakuNode is test scaffolding, so it moves to the test helpers along
with the port allocation. Neither is part of the package surface now.
2026-08-25 14:31:42 +01:00

1086 lines
33 KiB
Go

package kernel
import (
"context"
"errors"
"fmt"
"slices"
"testing"
"time"
"go.uber.org/zap"
"google.golang.org/protobuf/proto"
"github.com/cenkalti/backoff/v3"
"github.com/libp2p/go-libp2p/core/peer"
"github.com/logos-messaging/logos-delivery-go-bindings/pkg/kernel/common"
"github.com/logos-messaging/logos-delivery-go-bindings/pkg/kernel/pb"
"github.com/logos-messaging/logos-delivery-go-bindings/pkg/kernel/store"
ma "github.com/multiformats/go-multiaddr"
"github.com/stretchr/testify/require"
)
// In order to run this test, you must run an nwaku node
//
// Using Docker:
//
// IP_ADDRESS=$(hostname -I | awk '{print $1}');
// docker run \
// -p 61000:61000/tcp -p 8000:8000/udp -p 8646:8646/tcp harbor.status.im/wakuorg/nwaku:v0.33.0 \
// --discv5-discovery=true --cluster-id=16 --log-level=DEBUG --shard=64 --tcp-port=61000 \
// --nat=extip:${IP_ADDRESS} --discv5-udp-port=8000 --rest-address=0.0.0.0 --store --rest-port=8646
func TestBasicWaku(t *testing.T) {
requiresNode(t)
t.Skip("Skipping test as choosing this port will fail the CI")
extNodeRestPort := 8646
storeNodeInfo, err := GetNwakuInfo(nil, &extNodeRestPort)
require.NoError(t, err)
// ctx := context.Background()
nwakuConfig := common.WakuConfig{
Nodekey: "11d0dcea28e86f81937a3bd1163473c7fbc0a0db54fd72914849bc47bdf78710",
Relay: true,
LogLevel: "DEBUG",
DnsDiscoveryUrl: "enrtree://AMOJVZX4V6EXP7NTJPMAYJYST2QP6AJXYW76IU6VGJS7UVSNDYZG4@boot.prod.status.nodes.status.im",
DnsDiscovery: true,
Discv5Discovery: true,
Staticnodes: []string{storeNodeInfo.ListenAddresses[0]},
ClusterID: 16,
Shards: []uint16{64},
}
storeNodeMa, err := ma.NewMultiaddr(storeNodeInfo.ListenAddresses[0])
require.NoError(t, err)
w, err := NewFromWakuConfig(&nwakuConfig)
require.NoError(t, err)
require.NoError(t, w.Start())
enr, err := w.Debug().ENR()
require.NoError(t, err)
require.NotNil(t, enr)
options := func(b *backoff.ExponentialBackOff) {
b.MaxElapsedTime = 30 * time.Second
}
// Sanity check, not great, but it's probably helpful
err = RetryWithBackOff(func() error {
numConnected, err := w.Peers().NumConnected()
if err != nil {
return err
}
// Have to be connected to at least 3 nodes: the static node, the bootstrap node, and one discovered node
if numConnected > 2 {
return nil
}
return errors.New("no peers discovered")
}, options)
require.NoError(t, err)
// Get local store node address
storeNode, err := peer.AddrInfoFromString(storeNodeInfo.ListenAddresses[0])
require.NoError(t, err)
/*
w.node.Peers().Dial(ctx, storeNode.Addrs[0], "")
w.StorenodeCycle.SetStorenodeConfigProvider(newTestStorenodeConfigProvider(*storeNode))
*/
// Check that we are indeed connected to the store node
connectedStoreNodes, err := w.Peers().ByProtocol(store.StoreQueryID_v300)
require.NoError(t, err)
require.True(t, slices.Contains(connectedStoreNodes, storeNode.ID), "nwaku should be connected to the store node")
// Disconnect from the store node
err = w.Peers().Disconnect(storeNode.ID)
require.NoError(t, err)
// Check that we are indeed disconnected
connectedStoreNodes, err = w.Peers().ByProtocol(store.StoreQueryID_v300)
require.NoError(t, err)
isDisconnected := !slices.Contains(connectedStoreNodes, storeNode.ID)
require.True(t, isDisconnected, "nwaku should be disconnected from the store node")
// Re-connect
ctx, cancel := context.WithTimeout(context.Background(), requestTimeout)
defer cancel()
err = w.Peers().Connect(ctx, storeNodeMa)
require.NoError(t, err)
// Check that we are connected again
connectedStoreNodes, err = w.Peers().ByProtocol(store.StoreQueryID_v300)
require.NoError(t, err)
require.True(t, slices.Contains(connectedStoreNodes, storeNode.ID), "nwaku should be connected to the store node")
/* filter := &common.Filter{
PubsubTopic: w.cfg.DefaultShardPubsubTopic,
Messages: common.NewMemoryMessageStore(),
ContentTopics: common.NewTopicSetFromBytes([][]byte{{1, 2, 3, 4}}),
}
_, err = w.Subscribe(filter)
require.NoError(t, err)
msgTimestamp := w.timestamp()
contentTopic := maps.Keys(filter.ContentTopics)[0]
time.Sleep(2 * time.Second)
msgID, err := w.Send(w.cfg.DefaultShardPubsubTopic, &pb.WakuMessage{
Payload: []byte{1, 2, 3, 4, 5},
ContentTopic: contentTopic.ContentTopic(),
Version: proto.Uint32(0),
Timestamp: &msgTimestamp,
}, nil)
require.NoError(t, err)
require.NotEqual(t, msgID, "1")
time.Sleep(1 * time.Second)
messages := filter.Retrieve()
require.Len(t, messages, 1)
timestampInSeconds := msgTimestamp / int64(time.Second)
marginInSeconds := 20
options = func(b *backoff.ExponentialBackOff) {
b.MaxElapsedTime = 60 * time.Second
b.InitialInterval = 500 * time.Millisecond
}
err = RetryWithBackOff(func() error {
err := w.HistoryRetriever.Query(
context.Background(),
store.FilterCriteria{
ContentFilter: protocol.NewContentFilter(w.cfg.DefaultShardPubsubTopic, contentTopic.ContentTopic()),
TimeStart: proto.Int64((timestampInSeconds - int64(marginInSeconds)) * int64(time.Second)),
TimeEnd: proto.Int64((timestampInSeconds + int64(marginInSeconds)) * int64(time.Second)),
},
*storeNode,
10,
nil, false,
)
return err
// TODO-nwaku
if err != nil || envelopeCount == 0 {
// in case of failure extend timestamp margin up to 40secs
if marginInSeconds < 40 {
marginInSeconds += 5
}
return errors.New("no messages received from store node")
}
return nil
}, options)
require.NoError(t, err)
time.Sleep(10 * time.Second)
*/
require.NoError(t, w.Stop())
}
func TestPeerExchange(t *testing.T) {
requiresNode(t)
// start node that will be discovered by PeerExchange
discV5NodeWakuConfig := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: true,
ClusterID: 16,
Shards: []uint16{64},
PeerExchange: false,
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
discV5Node, err := NewFromWakuConfig(&discV5NodeWakuConfig)
require.NoError(t, err)
require.NoError(t, discV5Node.Start())
discV5NodePeerId, err := discV5Node.Debug().PeerID()
require.NoError(t, err)
discv5NodeEnr, err := discV5Node.Debug().ENR()
require.NoError(t, err)
// start node which serves as PeerExchange server
pxServerWakuConfig := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: true,
ClusterID: 16,
Shards: []uint16{64},
PeerExchange: true,
Discv5UdpPort: freeUDPPort(t),
Discv5BootstrapNodes: []string{discv5NodeEnr.String()},
TcpPort: freeTCPPort(t),
}
pxServerNode, err := NewFromWakuConfig(&pxServerWakuConfig)
require.NoError(t, err)
require.NoError(t, pxServerNode.Start())
// Adding an extra second to make sure PX cache is not empty
time.Sleep(2 * time.Second)
serverNodeMa, err := pxServerNode.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, serverNodeMa)
require.True(t, len(serverNodeMa) > 0)
// Sanity check, not great, but it's probably helpful
options := func(b *backoff.ExponentialBackOff) {
b.MaxElapsedTime = 30 * time.Second
}
// Check that pxServerNode has discV5Node in its Peer Store
err = RetryWithBackOff(func() error {
peers, err := pxServerNode.Peers().FromPeerStore()
if err != nil {
return err
}
if slices.Contains(peers, discV5NodePeerId) {
return nil
}
return errors.New("pxServer is missing the discv5 node in its peer store")
}, options)
require.NoError(t, err)
// start light node which uses PeerExchange to discover peers
pxClientWakuConfig := common.WakuConfig{
Relay: false,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: 16,
Shards: []uint16{64},
PeerExchange: true,
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
PeerExchangeNode: serverNodeMa[0].String(),
}
lightNode, err := NewFromWakuConfig(&pxClientWakuConfig)
require.NoError(t, err)
require.NoError(t, lightNode.Start())
pxServerPeerId, err := pxServerNode.Debug().PeerID()
require.NoError(t, err)
// Check that the light node discovered the discV5Node and has both nodes in its peer store
err = RetryWithBackOff(func() error {
peers, err := lightNode.Peers().FromPeerStore()
if err != nil {
return err
}
if slices.Contains(peers, discV5NodePeerId) && slices.Contains(peers, pxServerPeerId) {
return nil
}
return errors.New("lightnode is missing peers")
}, options)
require.NoError(t, err)
// Now perform the PX request manually to see if it also works
err = RetryWithBackOff(func() error {
numPeersReceived, err := lightNode.PeerExchange().Request(1)
if err != nil {
return err
}
if numPeersReceived == 1 {
return nil
}
return errors.New("Peer Exchange is not returning peers")
}, options)
require.NoError(t, err)
// Stop nodes
require.NoError(t, lightNode.Stop())
require.NoError(t, pxServerNode.Stop())
require.NoError(t, discV5Node.Stop())
}
func TestDnsDiscover(t *testing.T) {
requiresNode(t)
nameserver := "8.8.8.8"
nodeWakuConfig := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
ClusterID: 16,
Shards: []uint16{64},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
node, err := NewFromWakuConfig(&nodeWakuConfig)
require.NoError(t, err)
require.NoError(t, node.Start())
sampleEnrTree := "enrtree://AMOJVZX4V6EXP7NTJPMAYJYST2QP6AJXYW76IU6VGJS7UVSNDYZG4@boot.prod.status.nodes.status.im"
ctx, cancel := context.WithTimeout(context.TODO(), requestTimeout)
defer cancel()
res, err := node.DNSDiscovery().Resolve(ctx, sampleEnrTree, nameserver)
require.NoError(t, err)
require.True(t, len(res) > 1, "multiple nodes should be returned from the DNS Discovery query")
// Stop nodes
require.NoError(t, node.Stop())
}
func TestDial(t *testing.T) {
requiresNode(t)
// start node that will initiate the dial
dialerNodeWakuConfig := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: 16,
Shards: []uint16{64},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
dialerNode, err := NewFromWakuConfig(&dialerNodeWakuConfig)
require.NoError(t, err)
require.NoError(t, dialerNode.Start())
// start node that will receive the dial
receiverNodeWakuConfig := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: 16,
Shards: []uint16{64},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
receiverNode, err := NewFromWakuConfig(&receiverNodeWakuConfig)
require.NoError(t, err)
require.NoError(t, receiverNode.Start())
receiverMultiaddr, err := receiverNode.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, receiverMultiaddr)
require.True(t, len(receiverMultiaddr) > 0)
// Check that both nodes start with no connected peers
dialerPeerCount, err := dialerNode.Peers().NumConnected()
require.NoError(t, err)
require.True(t, dialerPeerCount == 0, "Dialer node should have no connected peers")
receiverPeerCount, err := receiverNode.Peers().NumConnected()
require.NoError(t, err)
require.True(t, receiverPeerCount == 0, "Receiver node should have no connected peers")
// Dial
ctx, cancel := context.WithTimeout(context.Background(), requestTimeout)
defer cancel()
err = dialerNode.Peers().Connect(ctx, receiverMultiaddr[0])
require.NoError(t, err)
time.Sleep(1 * time.Second)
// Check that both nodes now have one connected peer
dialerPeerCount, err = dialerNode.Peers().NumConnected()
require.NoError(t, err)
require.True(t, dialerPeerCount == 1, "Dialer node should have 1 peer")
receiverPeerCount, err = receiverNode.Peers().NumConnected()
require.NoError(t, err)
require.True(t, receiverPeerCount == 1, "Receiver node should have 1 peer")
// Stop nodes
require.NoError(t, dialerNode.Stop())
require.NoError(t, receiverNode.Stop())
}
func TestRelay(t *testing.T) {
requiresNode(t)
// start node that will send the message
senderNodeWakuConfig := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: 16,
Shards: []uint16{64},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
senderNode, err := NewFromWakuConfig(&senderNodeWakuConfig)
require.NoError(t, err)
require.NoError(t, senderNode.Start())
// start node that will receive the message
receiverNodeWakuConfig := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: 16,
Shards: []uint16{64},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
receiverNode, err := NewFromWakuConfig(&receiverNodeWakuConfig)
require.NoError(t, err)
require.NoError(t, receiverNode.Start())
receiverMultiaddr, err := receiverNode.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, receiverMultiaddr)
require.True(t, len(receiverMultiaddr) > 0)
// Dial so they become peers
ctx, cancel := context.WithTimeout(context.Background(), requestTimeout)
defer cancel()
err = senderNode.Peers().Connect(ctx, receiverMultiaddr[0])
require.NoError(t, err)
time.Sleep(1 * time.Second)
// Check that both nodes now have one connected peer
senderPeerCount, err := senderNode.Peers().NumConnected()
require.NoError(t, err)
require.True(t, senderPeerCount == 1, "Dialer node should have 1 peer")
receiverPeerCount, err := receiverNode.Peers().NumConnected()
require.NoError(t, err)
require.True(t, receiverPeerCount == 1, "Receiver node should have 1 peer")
message := &pb.WakuMessage{
Payload: []byte{1, 2, 3, 4, 5, 6},
ContentTopic: "test-content-topic",
Version: proto.Uint32(0),
Timestamp: proto.Int64(time.Now().UnixNano()),
}
// send message
pubsubTopic := FormatWakuRelayTopic(senderNodeWakuConfig.ClusterID, senderNodeWakuConfig.Shards[0])
ctx2, cancel2 := context.WithTimeout(context.Background(), requestTimeout)
defer cancel2()
_, _ = senderNode.Relay().Publish(ctx2, pubsubTopic, message)
// Wait to receive message
select {
case envelope := <-receiverNode.Messages():
require.NotNil(t, envelope, "Envelope should be received")
require.Equal(t, message.Payload, envelope.Message().Payload, "Received payload should match")
require.Equal(t, message.ContentTopic, envelope.Message().ContentTopic, "Content topic should match")
case <-time.After(10 * time.Second):
t.Fatal("Timeout: No message received within 10 seconds")
}
// Stop nodes
require.NoError(t, senderNode.Stop())
require.NoError(t, receiverNode.Stop())
}
func TestTopicHealth(t *testing.T) {
requiresNode(t)
clusterId := uint16(16)
shardId := uint16(64)
// start node1
wakuConfig1 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: clusterId,
Shards: []uint16{shardId},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
node1, err := NewFromWakuConfig(&wakuConfig1)
require.NoError(t, err)
require.NoError(t, node1.Start())
// start node2
wakuConfig2 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: clusterId,
Shards: []uint16{shardId},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
node2, err := NewFromWakuConfig(&wakuConfig2)
require.NoError(t, err)
require.NoError(t, node2.Start())
multiaddr2, err := node2.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, multiaddr2)
require.True(t, len(multiaddr2) > 0)
// node1 dials node2 so they become peers
ctx, cancel := context.WithTimeout(context.Background(), requestTimeout)
defer cancel()
err = node1.Peers().Connect(ctx, multiaddr2[0])
require.NoError(t, err)
time.Sleep(1 * time.Second)
// Check that both nodes now have one connected peer
peerCount1, err := node1.Peers().NumConnected()
require.NoError(t, err)
require.True(t, peerCount1 == 1, "node1 should have 1 peer")
peerCount2, err := node2.Peers().NumConnected()
require.NoError(t, err)
require.True(t, peerCount2 == 1, "node2 should have 1 peer")
// Wait to receive topic health update
select {
case topicHealth := <-node2.TopicHealthChanges():
require.NotNil(t, topicHealth, "topicHealth should be updated")
require.Equal(t, topicHealth.TopicHealth, "MinimallyHealthy", "Topic health should be MinimallyHealthy")
require.Equal(t, topicHealth.PubsubTopic, FormatWakuRelayTopic(clusterId, shardId), "PubsubTopic should match configured cluster and shard")
case <-time.After(10 * time.Second):
t.Fatal("Timeout: No topic health event received within 10 seconds")
}
// Stop nodes
require.NoError(t, node1.Stop())
require.NoError(t, node2.Stop())
}
func TestConnectionChange(t *testing.T) {
requiresNode(t)
clusterId := uint16(16)
shardId := uint16(64)
// start node1
wakuConfig1 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: clusterId,
Shards: []uint16{shardId},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
node1, err := NewFromWakuConfig(&wakuConfig1)
require.NoError(t, err)
require.NoError(t, node1.Start())
// start node2
wakuConfig2 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: clusterId,
Shards: []uint16{shardId},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
node2, err := NewFromWakuConfig(&wakuConfig2)
require.NoError(t, err)
require.NoError(t, node2.Start())
multiaddr2, err := node2.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, multiaddr2)
require.True(t, len(multiaddr2) > 0)
// node1 dials node2 so they become peers
ctx, cancel := context.WithTimeout(context.Background(), requestTimeout)
defer cancel()
err = node1.Peers().Connect(ctx, multiaddr2[0])
require.NoError(t, err)
time.Sleep(1 * time.Second)
// Check that both nodes now have one connected peer
peerCount1, err := node1.Peers().NumConnected()
require.NoError(t, err)
require.True(t, peerCount1 == 1, "node1 should have 1 peer")
peerCount2, err := node2.Peers().NumConnected()
require.NoError(t, err)
require.True(t, peerCount2 == 1, "node2 should have 1 peer")
peerId1, err := node1.Debug().PeerID()
require.NoError(t, err)
// Wait to receive connectionChange event
select {
case connectionChange := <-node2.ConnectionChanges():
require.NotNil(t, connectionChange, "connectionChange should be updated")
require.Equal(t, connectionChange.PeerEvent, "Joined", "connectionChange Joined event should be emitted")
require.Equal(t, connectionChange.PeerID, peerId1, "connectionChange event should contain node 1's peerId")
case <-time.After(10 * time.Second):
t.Fatal("Timeout: No connectionChange event received within 10 seconds")
}
// Disconnect from node1
err = node2.Peers().Disconnect(peerId1)
require.NoError(t, err)
// Wait to receive connectionChange event
select {
case connectionChange := <-node2.ConnectionChanges():
require.NotNil(t, connectionChange, "connectionChange should be updated")
require.Equal(t, connectionChange.PeerEvent, "Left", "connectionChange Left event should be emitted")
require.Equal(t, connectionChange.PeerID, peerId1, "connectionChange event should contain node 1's peerId")
case <-time.After(10 * time.Second):
t.Fatal("Timeout: No connectionChange event received within 10 seconds")
}
// Stop nodes
require.NoError(t, node1.Stop())
require.NoError(t, node2.Stop())
}
func TestStore(t *testing.T) {
requiresNode(t)
// start node that will send the message
senderNodeWakuConfig := common.WakuConfig{
Relay: true,
Store: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: 16,
Shards: []uint16{64},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
senderNode, err := NewFromWakuConfig(&senderNodeWakuConfig)
require.NoError(t, err)
require.NoError(t, senderNode.Start())
// start node that will receive the message
receiverNodeWakuConfig := common.WakuConfig{
Relay: true,
Store: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: 16,
Shards: []uint16{64},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
receiverNode, err := NewFromWakuConfig(&receiverNodeWakuConfig)
require.NoError(t, err)
require.NoError(t, receiverNode.Start())
receiverMultiaddr, err := receiverNode.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, receiverMultiaddr)
require.True(t, len(receiverMultiaddr) > 0)
// Dial so they become peers
ctx, cancel := context.WithTimeout(context.Background(), requestTimeout)
defer cancel()
err = senderNode.Peers().Connect(ctx, receiverMultiaddr[0])
require.NoError(t, err)
time.Sleep(1 * time.Second)
// Check that both nodes now have one connected peer
senderPeerCount, err := senderNode.Peers().NumConnected()
require.NoError(t, err)
require.True(t, senderPeerCount == 1, "Dialer node should have 1 peer")
receiverPeerCount, err := receiverNode.Peers().NumConnected()
require.NoError(t, err)
require.True(t, receiverPeerCount == 1, "Receiver node should have 1 peer")
// Send 8 messages
numMessages := 8
paginationLimit := 5
timeStart := proto.Int64(time.Now().UnixNano())
hashes := []common.MessageHash{}
pubsubTopic := FormatWakuRelayTopic(senderNodeWakuConfig.ClusterID, senderNodeWakuConfig.Shards[0])
for i := 0; i < numMessages; i++ {
message := &pb.WakuMessage{
Payload: []byte{byte(i)}, // Include message number in payload
ContentTopic: "test-content-topic",
Version: proto.Uint32(0),
Timestamp: proto.Int64(time.Now().UnixNano()),
}
ctx2, cancel2 := context.WithTimeout(context.Background(), requestTimeout)
defer cancel2()
hash, err := senderNode.Relay().Publish(ctx2, pubsubTopic, message)
require.NoError(t, err)
hashes = append(hashes, hash)
}
// Wait to receive all 8 messages
receivedCount := 0
receivedMessages := make(map[byte]bool)
// Use a timeout for the entire receive operation
timeoutChan := time.After(10 * time.Second)
for receivedCount < numMessages {
select {
case envelope := <-receiverNode.Messages():
require.NotNil(t, envelope, "Envelope should be received")
payload := envelope.Message().Payload
msgNum := payload[0]
// Check if we've already received this message number
if !receivedMessages[msgNum] {
receivedMessages[msgNum] = true
receivedCount++
}
require.Equal(t, "test-content-topic", envelope.Message().ContentTopic, "Content topic should match")
case <-timeoutChan:
t.Fatalf("Timeout: Only received %d messages out of 8 within 10 seconds", receivedCount)
}
}
// Verify we received all messages
for i := 0; i < numMessages; i++ {
require.True(t, receivedMessages[byte(i)], fmt.Sprintf("Message %d was not received", i))
}
// Now send store query
storeReq1 := common.StoreQueryRequest{
IncludeData: true,
ContentTopics: &[]string{"test-content-topic"},
PaginationLimit: proto.Uint64(uint64(paginationLimit)),
PaginationForward: true,
TimeStart: timeStart,
}
storeNodeAddrInfo, err := peer.AddrInfoFromString(receiverMultiaddr[0].String())
require.NoError(t, err)
ctx3, cancel3 := context.WithTimeout(context.Background(), requestTimeout)
defer cancel3()
res1, err := senderNode.Store().Query(ctx3, &storeReq1, *storeNodeAddrInfo)
require.NoError(t, err)
storedMessages1 := *res1.Messages
for i := 0; i < paginationLimit; i++ {
require.True(t, storedMessages1[i].MessageHash == hashes[i], fmt.Sprintf("Stored message does not match received message for index %d", i))
}
// Now let's query the second page
storeReq2 := common.StoreQueryRequest{
IncludeData: true,
ContentTopics: &[]string{"test-content-topic"},
PaginationLimit: proto.Uint64(uint64(paginationLimit)),
PaginationForward: true,
TimeStart: timeStart,
PaginationCursor: &res1.PaginationCursor,
}
ctx4, cancel4 := context.WithTimeout(context.Background(), requestTimeout)
defer cancel4()
res2, err := senderNode.Store().Query(ctx4, &storeReq2, *storeNodeAddrInfo)
require.NoError(t, err)
storedMessages2 := *res2.Messages
for i := 0; i < len(storedMessages2); i++ {
require.True(t, storedMessages2[i].MessageHash == hashes[i+paginationLimit], fmt.Sprintf("Stored message does not match received message for index %d", i))
}
// Now let's query for two specific message hashes
storeReq3 := common.StoreQueryRequest{
IncludeData: true,
ContentTopics: &[]string{"test-content-topic"},
MessageHashes: &[]common.MessageHash{hashes[0], hashes[2]},
}
ctx5, cancel5 := context.WithTimeout(context.Background(), requestTimeout)
defer cancel5()
res3, err := senderNode.Store().Query(ctx5, &storeReq3, *storeNodeAddrInfo)
require.NoError(t, err)
storedMessages3 := *res3.Messages
require.True(t, storedMessages3[0].MessageHash == hashes[0], "Stored message does not match queried message")
require.True(t, storedMessages3[1].MessageHash == hashes[2], "Stored message does not match queried message")
// Stop nodes
require.NoError(t, senderNode.Stop())
require.NoError(t, receiverNode.Stop())
}
func TestParallelPings(t *testing.T) {
requiresNode(t)
logger, err := zap.NewDevelopment()
require.NoError(t, err)
// start node that will initiate the dial
dialerNodeWakuConfig := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: 16,
Shards: []uint16{64},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
dialerNode, err := NewFromWakuConfig(&dialerNodeWakuConfig)
require.NoError(t, err)
require.NoError(t, dialerNode.Start())
receiverNodeWakuConfig1 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: 16,
Shards: []uint16{64},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
receiverNode1, err := NewFromWakuConfig(&receiverNodeWakuConfig1)
require.NoError(t, err)
require.NoError(t, receiverNode1.Start())
receiverMultiaddr1, err := receiverNode1.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, receiverMultiaddr1)
require.True(t, len(receiverMultiaddr1) > 0)
receiverNodeWakuConfig2 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: 16,
Shards: []uint16{64},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
receiverNode2, err := NewFromWakuConfig(&receiverNodeWakuConfig2)
require.NoError(t, err)
require.NoError(t, receiverNode2.Start())
receiverMultiaddr2, err := receiverNode2.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, receiverMultiaddr2)
require.True(t, len(receiverMultiaddr2) > 0)
receiverNodeWakuConfig3 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: 16,
Shards: []uint16{64},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
receiverNode3, err := NewFromWakuConfig(&receiverNodeWakuConfig3)
require.NoError(t, err)
require.NoError(t, receiverNode3.Start())
receiverMultiaddr3, err := receiverNode3.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, receiverMultiaddr3)
require.True(t, len(receiverMultiaddr3) > 0)
receiverNodes := []string{receiverMultiaddr1[0].String(), receiverMultiaddr2[0].String(), receiverMultiaddr3[0].String()}
// node.Peers().Ping(ctx, peerInfo)
for _, receiverNode := range receiverNodes {
addrInfo, err := peer.AddrInfoFromString(receiverNode)
require.NoError(t, err)
go func(peerInfo peer.AddrInfo) {
ctx, cancel := context.WithTimeout(context.Background(), requestTimeout)
defer cancel()
_, err := dialerNode.Peers().Ping(ctx, peerInfo)
if err != nil { // pinging storenodes might fail, but we don't care
logger.Warn("failed pinging node", zap.Stringer("peerId", addrInfo.ID), zap.Error(err))
}
}(*addrInfo)
}
options := func(b *backoff.ExponentialBackOff) {
b.MaxElapsedTime = 30 * time.Second
}
err = RetryWithBackOff(func() error {
dialerPeerCount, err := dialerNode.Peers().NumConnected()
if err != nil {
return err
}
if dialerPeerCount == 3 {
return nil
}
return fmt.Errorf("dialerNode should have 3 peers but it has %d", dialerPeerCount)
}, options)
require.NoError(t, err)
// Stop nodes
require.NoError(t, dialerNode.Stop())
}
func TestOnline(t *testing.T) {
requiresNode(t)
clusterId := uint16(16)
shardId := uint16(64)
// start node1
wakuConfig1 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: clusterId,
Shards: []uint16{shardId},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
node1, err := NewFromWakuConfig(&wakuConfig1)
require.NoError(t, err)
require.NoError(t, node1.Start())
// start node2
wakuConfig2 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: clusterId,
Shards: []uint16{shardId},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
node2, err := NewFromWakuConfig(&wakuConfig2)
require.NoError(t, err)
require.NoError(t, node2.Start())
multiaddr2, err := node2.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, multiaddr2)
require.True(t, len(multiaddr2) > 0)
// node1 dials node2 so they become peers
ctx, cancel := context.WithTimeout(context.Background(), requestTimeout)
defer cancel()
err = node1.Peers().Connect(ctx, multiaddr2[0])
require.NoError(t, err)
time.Sleep(1 * time.Second)
// Check that both nodes now have one connected peer
peerCount1, err := node1.Peers().NumConnected()
require.NoError(t, err)
require.True(t, peerCount1 == 1, "node1 should have 1 peer")
peerCount2, err := node2.Peers().NumConnected()
require.NoError(t, err)
require.True(t, peerCount2 == 1, "node2 should have 1 peer")
isOnline, err := node1.Debug().IsOnline()
require.NoError(t, err)
require.True(t, isOnline, "node1 should be online")
// Stop nodes
require.NoError(t, node1.Stop())
require.NoError(t, node2.Stop())
}
func TestDisconnectAllPeers(t *testing.T) {
requiresNode(t)
clusterId := uint16(16)
shardId := uint16(64)
// start node1
wakuConfig1 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: clusterId,
Shards: []uint16{shardId},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
node1, err := NewFromWakuConfig(&wakuConfig1)
require.NoError(t, err)
require.NoError(t, node1.Start())
defer node1.Stop()
// start node2
wakuConfig2 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: clusterId,
Shards: []uint16{shardId},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
node2, err := NewFromWakuConfig(&wakuConfig2)
require.NoError(t, err)
require.NoError(t, node2.Start())
defer node2.Stop()
multiaddr2, err := node2.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, multiaddr2)
require.True(t, len(multiaddr2) > 0)
// start node3
wakuConfig3 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: clusterId,
Shards: []uint16{shardId},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
node3, err := NewFromWakuConfig(&wakuConfig3)
require.NoError(t, err)
require.NoError(t, node3.Start())
defer node3.Stop()
multiaddr3, err := node3.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, multiaddr3)
require.True(t, len(multiaddr3) > 0)
// start node4
wakuConfig4 := common.WakuConfig{
Relay: true,
LogLevel: "DEBUG",
Discv5Discovery: false,
ClusterID: clusterId,
Shards: []uint16{shardId},
Discv5UdpPort: freeUDPPort(t),
TcpPort: freeTCPPort(t),
}
node4, err := NewFromWakuConfig(&wakuConfig4)
require.NoError(t, err)
require.NoError(t, node4.Start())
defer node4.Stop()
multiaddr4, err := node4.Debug().ListenAddresses()
require.NoError(t, err)
require.NotNil(t, multiaddr4)
require.True(t, len(multiaddr4) > 0)
to_dial := []ma.Multiaddr{multiaddr2[0], multiaddr3[0], multiaddr4[0]}
for _, addr := range to_dial {
ctx, cancel := context.WithTimeout(context.Background(), requestTimeout)
defer cancel()
err = node1.Peers().Connect(ctx, addr)
require.NoError(t, err)
}
time.Sleep(1 * time.Second)
peerCount1, err := node1.Peers().NumConnected()
require.NoError(t, err)
require.True(t, peerCount1 == 3, "node1 should have 3 peers")
err = node1.Peers().DisconnectAll()
require.NoError(t, err)
peerCount1, err = node1.Peers().NumConnected()
require.NoError(t, err)
require.True(t, peerCount1 == 0, "node1 should have 0 peers")
}