Skip to content
15 changes: 8 additions & 7 deletions pkg/p2p/libp2p/libp2p.go
Original file line number Diff line number Diff line change
Expand Up @@ -589,7 +589,7 @@ func (s *Service) handleIncoming(stream network.Stream) {
s.logger.Debug("stream handler: handshake: build remote multiaddrs", "peer_id", peerID, "error", err)
s.logger.Error(nil, "stream handler: handshake: build remote multiaddrs", "peer_id", peerID)
_ = handshakeStream.Reset()
_ = s.host.Network().ClosePeer(peerID)
_ = stream.Conn().Close()
return
}

Expand Down Expand Up @@ -621,7 +621,7 @@ func (s *Service) handleIncoming(stream network.Stream) {
s.logger.Debug("stream handler: handshake: handle failed", "peer_id", peerID, "error", err)
s.logger.Error(nil, "stream handler: handshake: handle failed", "peer_id", peerID)
_ = handshakeStream.Reset()
_ = s.host.Network().ClosePeer(peerID)
_ = stream.Conn().Close()
return
}

Expand Down Expand Up @@ -1051,7 +1051,7 @@ func (s *Service) Connect(ctx context.Context, addrs []ma.Multiaddr) (address *b
peerAddrs, err := s.peerMultiaddrs(ctx, peerID)
if err != nil {
_ = handshakeStream.Reset()
_ = s.host.Network().ClosePeer(peerID)
_ = stream.Conn().Close()
return nil, fmt.Errorf("build peer multiaddrs: %w", err)
}

Expand All @@ -1075,13 +1075,13 @@ func (s *Service) Connect(ctx context.Context, addrs []ma.Multiaddr) (address *b
)
if err != nil {
_ = handshakeStream.Reset()
_ = s.host.Network().ClosePeer(info.ID)
_ = stream.Conn().Close()
return nil, fmt.Errorf("handshake: %w", err)
}

if !i.FullNode {
_ = handshakeStream.Reset()
_ = s.host.Network().ClosePeer(info.ID)
_ = stream.Conn().Close()
return nil, p2p.ErrDialLightNode
}

Expand All @@ -1092,7 +1092,7 @@ func (s *Service) Connect(ctx context.Context, addrs []ma.Multiaddr) (address *b
s.logger.Debug("blocklisting: exists failed", "peer_id", info.ID, "error", err)
s.logger.Error(nil, "internal error while connecting with peer", "peer_id", info.ID)
_ = handshakeStream.Reset()
_ = s.host.Network().ClosePeer(info.ID)
_ = stream.Conn().Close()
return nil, err
}

Expand All @@ -1105,7 +1105,8 @@ func (s *Service) Connect(ctx context.Context, addrs []ma.Multiaddr) (address *b

if exists := s.peers.addIfNotExists(stream.Conn(), overlay, i.FullNode); exists {
if err := handshakeStream.FullClose(); err != nil {
_ = s.Disconnect(overlay, "failed closing handshake stream after connect")
// Only close the new (duplicate) connection; keep existing healthy sessions intact.
_ = stream.Conn().Close()
return nil, fmt.Errorf("peer exists, full close: %w", err)
}

Expand Down
Loading