Skip to content

Commit 953ec44

Browse files
Merge remote-tracking branch 'origin/master' into feat/empty-underlay-support
2 parents fa87414 + 6ddf9b4 commit 953ec44

7 files changed

Lines changed: 344 additions & 47 deletions

File tree

pkg/hive/hive.go

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -50,7 +50,7 @@ var (
5050
)
5151

5252
type Service struct {
53-
streamer p2p.Streamer
53+
streamer p2p.Bee260CompatibilityStreamer
5454
addressBook addressbook.GetPutter
5555
addPeersHandler func(...swarm.Address)
5656
networkID uint64
@@ -67,7 +67,7 @@ type Service struct {
6767
overlay swarm.Address
6868
}
6969

70-
func New(streamer p2p.Streamer, addressbook addressbook.GetPutter, networkID uint64, bootnode bool, allowPrivateCIDRs bool, overlay swarm.Address, logger log.Logger) *Service {
70+
func New(streamer p2p.Bee260CompatibilityStreamer, addressbook addressbook.GetPutter, networkID uint64, bootnode bool, allowPrivateCIDRs bool, overlay swarm.Address, logger log.Logger) *Service {
7171
svc := &Service{
7272
streamer: streamer,
7373
logger: logger.WithName(loggerName).Register(),
@@ -196,6 +196,8 @@ func (s *Service) sendPeers(ctx context.Context, peer swarm.Address, peers []swa
196196
continue
197197
}
198198

199+
advertisableUnderlays = p2p.FilterBee260CompatibleUnderlays(s.streamer.IsBee260(peer), advertisableUnderlays)
200+
199201
peersRequest.Peers = append(peersRequest.Peers, &pb.BzzAddress{
200202
Overlay: addr.Overlay.Bytes(),
201203
Underlay: bzz.SerializeUnderlays(advertisableUnderlays),

pkg/p2p/libp2p/internal/handshake/handshake.go

Lines changed: 4 additions & 19 deletions
Original file line numberDiff line numberDiff line change
@@ -153,7 +153,7 @@ func (s *Service) Handshake(ctx context.Context, stream p2p.Stream, peerMultiadd
153153

154154
w, r := protobuf.NewWriterAndReader(stream)
155155

156-
peerMultiaddrs = filterBee260CompatibleUnderlays(o.bee260compatibility, peerMultiaddrs)
156+
peerMultiaddrs = p2p.FilterBee260CompatibleUnderlays(o.bee260compatibility, peerMultiaddrs)
157157

158158
if err := w.WriteMsgWithContext(ctx, &pb.Syn{
159159
ObservedUnderlay: bzz.SerializeUnderlays(peerMultiaddrs),
@@ -208,7 +208,7 @@ func (s *Service) Handshake(ctx context.Context, stream p2p.Stream, peerMultiadd
208208
return a.Equal(b)
209209
})
210210

211-
advertisableUnderlays = filterBee260CompatibleUnderlays(o.bee260compatibility, advertisableUnderlays)
211+
advertisableUnderlays = p2p.FilterBee260CompatibleUnderlays(o.bee260compatibility, advertisableUnderlays)
212212

213213
bzzAddress, err := bzz.NewAddress(s.signer, advertisableUnderlays, s.overlay, s.networkID, s.nonce)
214214
if err != nil {
@@ -306,7 +306,7 @@ func (s *Service) Handle(ctx context.Context, stream p2p.Stream, peerMultiaddrs
306306
return a.Equal(b)
307307
})
308308

309-
advertisableUnderlays = filterBee260CompatibleUnderlays(o.bee260compatibility, advertisableUnderlays)
309+
advertisableUnderlays = p2p.FilterBee260CompatibleUnderlays(o.bee260compatibility, advertisableUnderlays)
310310

311311
bzzAddress, err := bzz.NewAddress(s.signer, advertisableUnderlays, s.overlay, s.networkID, s.nonce)
312312
if err != nil {
@@ -315,7 +315,7 @@ func (s *Service) Handle(ctx context.Context, stream p2p.Stream, peerMultiaddrs
315315

316316
welcomeMessage := s.GetWelcomeMessage()
317317

318-
peerMultiaddrs = filterBee260CompatibleUnderlays(o.bee260compatibility, peerMultiaddrs)
318+
peerMultiaddrs = p2p.FilterBee260CompatibleUnderlays(o.bee260compatibility, peerMultiaddrs)
319319

320320
if err := w.WriteMsgWithContext(ctx, &pb.SynAck{
321321
Syn: &pb.Syn{
@@ -395,18 +395,3 @@ func (s *Service) parseCheckAck(ack *pb.Ack) (*bzz.Address, error) {
395395

396396
return bzzAddress, nil
397397
}
398-
399-
// filterBee260CompatibleUnderlays select a single underlay to pass if
400-
// bee260compatibility is true. Otherwise it passes the unmodified underlays
401-
// slice. This function can be safely removed when bee version 2.6.0 is
402-
// deprecated.
403-
func filterBee260CompatibleUnderlays(bee260compatibility bool, underlays []ma.Multiaddr) []ma.Multiaddr {
404-
if !bee260compatibility {
405-
return underlays
406-
}
407-
underlay := bzz.SelectBestAdvertisedAddress(underlays, nil)
408-
if underlay == nil {
409-
return underlays
410-
}
411-
return []ma.Multiaddr{underlay}
412-
}

pkg/p2p/libp2p/libp2p.go

Lines changed: 41 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -646,7 +646,7 @@ func (s *Service) handleIncoming(stream network.Stream) {
646646
peerID := stream.Conn().RemotePeer()
647647
handshakeStream := newStream(stream, s.metrics)
648648

649-
peerMultiaddrs, err := s.peerMultiaddrs(s.ctx, peerID)
649+
peerMultiaddrs, err := s.peerMultiaddrs(s.ctx, stream.Conn().RemoteMultiaddr(), peerID)
650650
if err != nil {
651651
s.logger.Debug("stream handler: handshake: build remote multiaddrs", "peer_id", peerID, "error", err)
652652
s.logger.Error(nil, "stream handler: handshake: build remote multiaddrs", "peer_id", peerID)
@@ -655,11 +655,13 @@ func (s *Service) handleIncoming(stream network.Stream) {
655655
return
656656
}
657657

658+
bee260Compat := s.bee260BackwardCompatibility(peerID)
659+
658660
i, err := s.handshakeService.Handle(
659661
s.ctx,
660662
handshakeStream,
661663
peerMultiaddrs,
662-
handshake.WithBee260Compatibility(s.bee260BackwardCompatibility(peerID)),
664+
handshake.WithBee260Compatibility(bee260Compat),
663665
)
664666
if err != nil {
665667
s.logger.Debug("stream handler: handshake: handle failed", "peer_id", peerID, "error", err)
@@ -1070,18 +1072,20 @@ func (s *Service) Connect(ctx context.Context, addrs []ma.Multiaddr) (address *b
10701072

10711073
handshakeStream := newStream(stream, s.metrics)
10721074

1073-
peerMultiaddrs, err := s.peerMultiaddrs(ctx, peerID)
1075+
peerMultiaddrs, err := s.peerMultiaddrs(ctx, stream.Conn().RemoteMultiaddr(), peerID)
10741076
if err != nil {
10751077
_ = handshakeStream.Reset()
10761078
_ = s.host.Network().ClosePeer(peerID)
10771079
return nil, fmt.Errorf("build peer multiaddrs: %w", err)
10781080
}
10791081

1082+
bee260Compat := s.bee260BackwardCompatibility(peerID)
1083+
10801084
i, err := s.handshakeService.Handshake(
10811085
s.ctx,
10821086
handshakeStream,
10831087
peerMultiaddrs,
1084-
handshake.WithBee260Compatibility(s.bee260BackwardCompatibility(peerID)),
1088+
handshake.WithBee260Compatibility(bee260Compat),
10851089
)
10861090
if err != nil {
10871091
_ = handshakeStream.Reset()
@@ -1471,18 +1475,37 @@ func (s *Service) determineCurrentNetworkStatus(err error) error {
14711475
}
14721476

14731477
// peerMultiaddrs builds full multiaddresses for a peer given information from
1474-
// libp2p host peerstore and falling back to the remote address from the
1475-
// connection.
1476-
func (s *Service) peerMultiaddrs(ctx context.Context, peerID libp2ppeer.ID) ([]ma.Multiaddr, error) {
1478+
// the libp2p host peerstore. If the peerstore doesn't have addresses yet,
1479+
// it falls back to using the remote address from the active connection.
1480+
func (s *Service) peerMultiaddrs(ctx context.Context, remoteAddr ma.Multiaddr, peerID libp2ppeer.ID) ([]ma.Multiaddr, error) {
14771481
waitPeersCtx, cancel := context.WithTimeout(ctx, peerstoreWaitAddrsTimeout)
14781482
defer cancel()
14791483

1480-
return buildFullMAs(waitPeerAddrs(waitPeersCtx, s.host.Peerstore(), peerID), peerID)
1484+
mas := waitPeerAddrs(waitPeersCtx, s.host.Peerstore(), peerID)
1485+
if len(mas) == 0 && remoteAddr != nil {
1486+
mas = []ma.Multiaddr{remoteAddr}
1487+
}
1488+
1489+
return buildFullMAs(mas, peerID)
1490+
}
1491+
1492+
// IsBee260 implements p2p.Bee260CompatibilityStreamer interface.
1493+
// It checks if a peer is running Bee version older than 2.7.0.
1494+
func (s *Service) IsBee260(overlay swarm.Address) bool {
1495+
peerID, found := s.peers.peerID(overlay)
1496+
if !found {
1497+
return false
1498+
}
1499+
return s.bee260BackwardCompatibility(peerID)
14811500
}
14821501

14831502
var version270 = *semver.Must(semver.NewVersion("2.7.0"))
14841503

14851504
func (s *Service) bee260BackwardCompatibility(peerID libp2ppeer.ID) bool {
1505+
if compat, found := s.peers.bee260(peerID); found {
1506+
return compat
1507+
}
1508+
14861509
userAgent := s.peerUserAgent(s.ctx, peerID)
14871510
p := strings.SplitN(userAgent, " ", 2)
14881511
if len(p) != 2 {
@@ -1493,7 +1516,16 @@ func (s *Service) bee260BackwardCompatibility(peerID libp2ppeer.ID) bool {
14931516
if err != nil {
14941517
return false
14951518
}
1496-
return v.LessThan(version270)
1519+
1520+
// Compare major.minor.patch only (ignore pre-release)
1521+
// This way 2.7.0-rc12 is treated as >= 2.7.0
1522+
vCore, err := semver.NewVersion(fmt.Sprintf("%d.%d.%d", v.Major, v.Minor, v.Patch))
1523+
if err != nil {
1524+
return false
1525+
}
1526+
result := vCore.LessThan(version270)
1527+
s.peers.setBee260(peerID, result)
1528+
return result
14971529
}
14981530

14991531
// appendSpace adds a leading space character if the string is not empty.

pkg/p2p/libp2p/peer.go

Lines changed: 34 additions & 17 deletions
Original file line numberDiff line numberDiff line change
@@ -18,12 +18,13 @@ import (
1818
)
1919

2020
type peerRegistry struct {
21-
underlays map[string]libp2ppeer.ID // map overlay address to underlay peer id
22-
overlays map[libp2ppeer.ID]swarm.Address // map underlay peer id to overlay address
23-
full map[libp2ppeer.ID]bool // map to track whether a node is full or light node (true=full)
24-
connections map[libp2ppeer.ID]map[network.Conn]struct{} // list of connections for safe removal on Disconnect notification
25-
streams map[libp2ppeer.ID]map[network.Stream]context.CancelFunc
26-
mu sync.RWMutex
21+
overlayToPeerID map[string]libp2ppeer.ID // map overlay address to underlay peer id
22+
overlays map[libp2ppeer.ID]swarm.Address // map underlay peer id to overlay address
23+
full map[libp2ppeer.ID]bool // map to track whether a node is full or light node (true=full)
24+
bee260Compatibility map[libp2ppeer.ID]bool // map to track bee260 backward compatibility
25+
connections map[libp2ppeer.ID]map[network.Conn]struct{} // list of connections for safe removal on Disconnect notification
26+
streams map[libp2ppeer.ID]map[network.Stream]context.CancelFunc
27+
mu sync.RWMutex
2728

2829
//nolint:misspell
2930
disconnecter disconnecter // peerRegistry notifies libp2p on peer disconnection
@@ -36,11 +37,12 @@ type disconnecter interface {
3637

3738
func newPeerRegistry() *peerRegistry {
3839
return &peerRegistry{
39-
underlays: make(map[string]libp2ppeer.ID),
40-
overlays: make(map[libp2ppeer.ID]swarm.Address),
41-
full: make(map[libp2ppeer.ID]bool),
42-
connections: make(map[libp2ppeer.ID]map[network.Conn]struct{}),
43-
streams: make(map[libp2ppeer.ID]map[network.Stream]context.CancelFunc),
40+
overlayToPeerID: make(map[string]libp2ppeer.ID),
41+
overlays: make(map[libp2ppeer.ID]swarm.Address),
42+
full: make(map[libp2ppeer.ID]bool),
43+
bee260Compatibility: make(map[libp2ppeer.ID]bool),
44+
connections: make(map[libp2ppeer.ID]map[network.Conn]struct{}),
45+
streams: make(map[libp2ppeer.ID]map[network.Stream]context.CancelFunc),
4446

4547
Notifiee: new(network.NoopNotifiee),
4648
}
@@ -75,12 +77,13 @@ func (r *peerRegistry) Disconnected(_ network.Network, c network.Conn) {
7577
delete(r.connections, peerID)
7678
overlay := r.overlays[peerID]
7779
delete(r.overlays, peerID)
78-
delete(r.underlays, overlay.ByteString())
80+
delete(r.overlayToPeerID, overlay.ByteString())
7981
for _, cancel := range r.streams[peerID] {
8082
cancel()
8183
}
8284
delete(r.streams, peerID)
8385
delete(r.full, peerID)
86+
delete(r.bee260Compatibility, peerID)
8487
r.mu.Unlock()
8588
r.disconnecter.disconnected(overlay)
8689

@@ -143,12 +146,12 @@ func (r *peerRegistry) addIfNotExists(c network.Conn, overlay swarm.Address, ful
143146
// this is solving a case of multiple underlying libp2p connections for the same peer
144147
r.connections[peerID][c] = struct{}{}
145148

146-
if _, exists := r.underlays[overlay.ByteString()]; exists {
149+
if _, exists := r.overlayToPeerID[overlay.ByteString()]; exists {
147150
return true
148151
}
149152

150153
r.streams[peerID] = make(map[network.Stream]context.CancelFunc)
151-
r.underlays[overlay.ByteString()] = peerID
154+
r.overlayToPeerID[overlay.ByteString()] = peerID
152155
r.overlays[peerID] = overlay
153156
r.full[peerID] = full
154157
return false
@@ -157,7 +160,7 @@ func (r *peerRegistry) addIfNotExists(c network.Conn, overlay swarm.Address, ful
157160

158161
func (r *peerRegistry) peerID(overlay swarm.Address) (peerID libp2ppeer.ID, found bool) {
159162
r.mu.RLock()
160-
peerID, found = r.underlays[overlay.ByteString()]
163+
peerID, found = r.overlayToPeerID[overlay.ByteString()]
161164
r.mu.RUnlock()
162165
return peerID, found
163166
}
@@ -176,6 +179,19 @@ func (r *peerRegistry) fullnode(peerID libp2ppeer.ID) (bool, bool) {
176179
return full, found
177180
}
178181

182+
func (r *peerRegistry) bee260(peerID libp2ppeer.ID) (compat, found bool) {
183+
r.mu.RLock()
184+
defer r.mu.RUnlock()
185+
compat, found = r.bee260Compatibility[peerID]
186+
return compat, found
187+
}
188+
189+
func (r *peerRegistry) setBee260(peerID libp2ppeer.ID, compat bool) {
190+
r.mu.Lock()
191+
defer r.mu.Unlock()
192+
r.bee260Compatibility[peerID] = compat
193+
}
194+
179195
func (r *peerRegistry) isConnected(peerID libp2ppeer.ID, remoteAddr ma.Multiaddr) (swarm.Address, bool) {
180196
if remoteAddr == nil {
181197
return swarm.ZeroAddress, false
@@ -207,16 +223,17 @@ func (r *peerRegistry) isConnected(peerID libp2ppeer.ID, remoteAddr ma.Multiaddr
207223

208224
func (r *peerRegistry) remove(overlay swarm.Address) (found, full bool, peerID libp2ppeer.ID) {
209225
r.mu.Lock()
210-
peerID, found = r.underlays[overlay.ByteString()]
226+
peerID, found = r.overlayToPeerID[overlay.ByteString()]
211227
delete(r.overlays, peerID)
212-
delete(r.underlays, overlay.ByteString())
228+
delete(r.overlayToPeerID, overlay.ByteString())
213229
delete(r.connections, peerID)
214230
for _, cancel := range r.streams[peerID] {
215231
cancel()
216232
}
217233
delete(r.streams, peerID)
218234
full = r.full[peerID]
219235
delete(r.full, peerID)
236+
delete(r.bee260Compatibility, peerID)
220237
r.mu.Unlock()
221238

222239
return found, full, peerID

0 commit comments

Comments
 (0)