Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

p2p: support EnableGossipService in p2p streams #6073

Merged
Merged
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion network/p2p/p2p.go
Original file line number Diff line number Diff line change
Expand Up @@ -159,7 +159,7 @@
// MakeService creates a P2P service instance
func MakeService(ctx context.Context, log logging.Logger, cfg config.Local, h host.Host, listenAddr string, wsStreamHandler StreamHandler, bootstrapPeers []*peer.AddrInfo) (*serviceImpl, error) {

sm := makeStreamManager(ctx, log, h, wsStreamHandler)
sm := makeStreamManager(ctx, log, h, wsStreamHandler, cfg.EnableGossipService)

Check warning on line 162 in network/p2p/p2p.go

View check run for this annotation

Codecov / codecov/patch

network/p2p/p2p.go#L162

Added line #L162 was not covered by tests
h.Network().Notify(sm)
h.SetStreamHandler(AlgorandWsProtocol, sm.streamHandler)

Expand Down
58 changes: 24 additions & 34 deletions network/p2p/streams.go
Original file line number Diff line number Diff line change
Expand Up @@ -30,10 +30,11 @@

// streamManager implements network.Notifiee to create and manage streams for use with non-gossipsub protocols.
type streamManager struct {
ctx context.Context
log logging.Logger
host host.Host
handler StreamHandler
ctx context.Context
log logging.Logger
host host.Host
handler StreamHandler
allowIncomingGossip bool

streams map[peer.ID]network.Stream
streamsLock deadlock.Mutex
Expand All @@ -42,18 +43,25 @@
// StreamHandler is called when a new bidirectional stream for a given protocol and peer is opened.
type StreamHandler func(ctx context.Context, pid peer.ID, s network.Stream, incoming bool)

func makeStreamManager(ctx context.Context, log logging.Logger, h host.Host, handler StreamHandler) *streamManager {
func makeStreamManager(ctx context.Context, log logging.Logger, h host.Host, handler StreamHandler, allowIncomingGossip bool) *streamManager {

Check warning on line 46 in network/p2p/streams.go

View check run for this annotation

Codecov / codecov/patch

network/p2p/streams.go#L46

Added line #L46 was not covered by tests
return &streamManager{
ctx: ctx,
log: log,
host: h,
handler: handler,
streams: make(map[peer.ID]network.Stream),
ctx: ctx,
log: log,
host: h,
handler: handler,
allowIncomingGossip: allowIncomingGossip,
streams: make(map[peer.ID]network.Stream),

Check warning on line 53 in network/p2p/streams.go

View check run for this annotation

Codecov / codecov/patch

network/p2p/streams.go#L48-L53

Added lines #L48 - L53 were not covered by tests
}
}

// streamHandler is called by libp2p when a new stream is accepted
func (n *streamManager) streamHandler(stream network.Stream) {
if stream.Conn().Stat().Direction == network.DirInbound && !n.allowIncomingGossip {
jasonpaulos marked this conversation as resolved.
Show resolved Hide resolved
n.log.Debugf("rejecting stream from incoming connection from %s", stream.Conn().RemotePeer().String())
stream.Close()
return

Check warning on line 62 in network/p2p/streams.go

View check run for this annotation

Codecov / codecov/patch

network/p2p/streams.go#L59-L62

Added lines #L59 - L62 were not covered by tests
}

n.streamsLock.Lock()
defer n.streamsLock.Unlock()

Expand All @@ -74,15 +82,7 @@
}
n.streams[stream.Conn().RemotePeer()] = stream

// streamHandler is supposed to be called for accepted streams, so we expect incoming here
incoming := stream.Conn().Stat().Direction == network.DirInbound
if !incoming {
if stream.Stat().Direction == network.DirUnknown {
n.log.Warnf("Unknown direction for a steam %s to/from %s", stream.ID(), remotePeer)
} else {
n.log.Warnf("Unexpected outgoing stream in streamHandler for connection %s (%s): %s vs %s stream", stream.Conn().ID(), remotePeer, stream.Conn().Stat().Direction, stream.Stat().Direction.String())
}
}
n.handler(n.ctx, remotePeer, stream, incoming)
return
}
Expand All @@ -92,20 +92,18 @@
}
// no old stream
n.streams[stream.Conn().RemotePeer()] = stream
// streamHandler is supposed to be called for accepted streams, so we expect incoming here
incoming := stream.Conn().Stat().Direction == network.DirInbound
if !incoming {
if stream.Stat().Direction == network.DirUnknown {
n.log.Warnf("streamHandler: unknown direction for a steam %s to/from %s", stream.ID(), remotePeer)
} else {
n.log.Warnf("Unexpected outgoing stream in streamHandler for connection %s (%s): %s vs %s stream", stream.Conn().ID(), remotePeer, stream.Conn().Stat().Direction, stream.Stat().Direction.String())
}
}
n.handler(n.ctx, remotePeer, stream, incoming)
}

// Connected is called when a connection is opened
// for both incoming (listener -> addConn) and outgoing (dialer -> addConn) connections.
func (n *streamManager) Connected(net network.Network, conn network.Conn) {
if conn.Stat().Direction == network.DirInbound && !n.allowIncomingGossip {
gmalouf marked this conversation as resolved.
Show resolved Hide resolved
n.log.Debugf("ignoring incoming connection from %s", conn.RemotePeer().String())
return

Check warning on line 104 in network/p2p/streams.go

View check run for this annotation

Codecov / codecov/patch

network/p2p/streams.go#L102-L104

Added lines #L102 - L104 were not covered by tests
jasonpaulos marked this conversation as resolved.
Show resolved Hide resolved
}

remotePeer := conn.RemotePeer()
localPeer := n.host.ID()

Expand Down Expand Up @@ -138,15 +136,7 @@
needUnlock = false
n.streamsLock.Unlock()

// a new stream created above, expected direction is outbound
incoming := stream.Conn().Stat().Direction == network.DirInbound
if incoming {
n.log.Warnf("Unexpected incoming stream in streamHandler for connection %s (%s): %s vs %s stream", stream.Conn().ID(), remotePeer, stream.Conn().Stat().Direction, stream.Stat().Direction.String())
} else {
if stream.Stat().Direction == network.DirUnknown {
n.log.Warnf("Connected: unknown direction for a steam %s to/from %s", stream.ID(), remotePeer)
}
}
n.handler(n.ctx, remotePeer, stream, incoming)
}

Expand Down
132 changes: 132 additions & 0 deletions network/p2pNetwork_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -1154,3 +1154,135 @@ func TestMergeP2PAddrInfoResolvedAddresses(t *testing.T) {
})
}
}

// TestP2PEnableGossipService_NodeDisable ensures that a node with EnableGossipService=false
// still can participate in the network by sending and receiving messages.
func TestP2PEnableGossipService_NodeDisable(t *testing.T) {
partitiontest.PartitionTest(t)

log := logging.TestingLog(t)

// prepare configs
cfg := config.GetDefaultLocal()
cfg.DNSBootstrapID = "" // disable DNS lookups since the test uses phonebook addresses

relayCfg := cfg
relayCfg.NetAddress = "127.0.0.1:0"

nodeCfg := cfg
nodeCfg.EnableGossipService = false
nodeCfg2 := nodeCfg
algorandskiy marked this conversation as resolved.
Show resolved Hide resolved
nodeCfg2.NetAddress = "127.0.0.1:0"

tests := []struct {
name string
relayCfg config.Local
nodeCfg config.Local
}{
{"non-listening-node", relayCfg, nodeCfg},
{"listening-node", relayCfg, nodeCfg2},
}
for _, test := range tests {
t.Run(test.name, func(t *testing.T) {
relayCfg := test.relayCfg
netA, err := NewP2PNetwork(log, relayCfg, "", nil, genesisID, config.Devtestnet, &nopeNodeInfo{}, nil)
require.NoError(t, err)
netA.Start()
defer netA.Stop()

peerInfoA := netA.service.AddrInfo()
addrsA, err := peer.AddrInfoToP2pAddrs(&peerInfoA)
require.NoError(t, err)
require.NotZero(t, addrsA[0])
multiAddrStr := addrsA[0].String()
phoneBookAddresses := []string{multiAddrStr}

// start netB with gossip service disabled
nodeCfg := test.nodeCfg
netB, err := NewP2PNetwork(log, nodeCfg, "", phoneBookAddresses, genesisID, config.Devtestnet, &nopeNodeInfo{}, nil)
require.NoError(t, err)
netB.Start()
defer netB.Stop()

require.Eventually(t, func() bool {
return netA.hasPeers() && netB.hasPeers()
}, 1*time.Second, 50*time.Millisecond)

testTag := protocol.AgreementVoteTag

var handlerCountA atomic.Uint32
passThroughHandlerA := []TaggedMessageHandler{
{Tag: testTag, MessageHandler: HandlerFunc(func(msg IncomingMessage) OutgoingMessage {
handlerCountA.Add(1)
return OutgoingMessage{Action: Broadcast}
})},
}
var handlerCountB atomic.Uint32
passThroughHandlerB := []TaggedMessageHandler{
{Tag: testTag, MessageHandler: HandlerFunc(func(msg IncomingMessage) OutgoingMessage {
handlerCountB.Add(1)
return OutgoingMessage{Action: Broadcast}
})},
}
netA.RegisterHandlers(passThroughHandlerA)
netB.RegisterHandlers(passThroughHandlerB)

// send messages from B and confirm that they get received by C (via A)
jasonpaulos marked this conversation as resolved.
Show resolved Hide resolved
for i := 0; i < 10; i++ {
err = netA.Broadcast(context.Background(), testTag, []byte(fmt.Sprintf("hello from A %d", i)), false, nil)
require.NoError(t, err)
err = netB.Broadcast(context.Background(), testTag, []byte(fmt.Sprintf("hello from B %d", i)), false, nil)
require.NoError(t, err)
}

require.Eventually(
t,
func() bool {
return handlerCountA.Load() == 10 && handlerCountB.Load() == 10
},
2*time.Second,
50*time.Millisecond,
)
})
}
}

// TestP2PEnableGossipService_BothDisable checks if both relay and node have EnableGossipService=false
// they do not connect to each other.
func TestP2PEnableGossipService_BothDisable(t *testing.T) {
partitiontest.PartitionTest(t)

log := logging.TestingLog(t)

// prepare configs
cfg := config.GetDefaultLocal()
cfg.DNSBootstrapID = "" // disable DNS lookups since the test uses phonebook addresses
cfg.EnableGossipService = false // disable gossip service by default

relayCfg := cfg
relayCfg.NetAddress = "127.0.0.1:0"

netA, err := NewP2PNetwork(log, relayCfg, "", nil, genesisID, config.Devtestnet, &nopeNodeInfo{}, nil)
require.NoError(t, err)
netA.Start()
defer netA.Stop()

peerInfoA := netA.service.AddrInfo()
addrsA, err := peer.AddrInfoToP2pAddrs(&peerInfoA)
require.NoError(t, err)
require.NotZero(t, addrsA[0])
multiAddrStr := addrsA[0].String()
phoneBookAddresses := []string{multiAddrStr}

nodeCfg := cfg
nodeCfg.NetAddress = ""

netB, err := NewP2PNetwork(log, nodeCfg, "", phoneBookAddresses, genesisID, config.Devtestnet, &nopeNodeInfo{}, nil)
require.NoError(t, err)
netB.Start()
defer netB.Stop()

require.Eventually(t, func() bool {
return !netA.hasPeers() && !netB.hasPeers()
}, 1*time.Second, 50*time.Millisecond)
jasonpaulos marked this conversation as resolved.
Show resolved Hide resolved
}
Loading