diff --git a/internal/p2p/channel.go b/internal/p2p/channel.go index d3d7d104f..d6763543a 100644 --- a/internal/p2p/channel.go +++ b/internal/p2p/channel.go @@ -46,6 +46,7 @@ type Wrapper interface { type PeerError struct { NodeID types.NodeID Err error + Fatal bool } func (pe PeerError) Error() string { return fmt.Sprintf("peer=%q: %s", pe.NodeID, pe.Err.Error()) } diff --git a/internal/p2p/peermanager.go b/internal/p2p/peermanager.go index 165b00e61..7391de4ea 100644 --- a/internal/p2p/peermanager.go +++ b/internal/p2p/peermanager.go @@ -430,6 +430,13 @@ func (m *PeerManager) PeerRatio() float64 { return float64(m.store.Size()) / float64(m.options.MaxPeers) } +func (m *PeerManager) HasMaxPeerCapacity() bool { + m.mtx.Lock() + defer m.mtx.Unlock() + + return len(m.connected) >= int(m.options.MaxConnected) +} + // DialNext finds an appropriate peer address to dial, and marks it as dialing. // If no peer is found, or all connection slots are full, it blocks until one // becomes available. The caller must call Dialed() or DialFailed() for the diff --git a/internal/p2p/router.go b/internal/p2p/router.go index a9a01f3c7..267d55a96 100644 --- a/internal/p2p/router.go +++ b/internal/p2p/router.go @@ -396,9 +396,21 @@ func (r *Router) routeChannel( return } - r.logger.Error("peer error, evicting", "peer", peerError.NodeID, "err", peerError.Err) + shouldEvict := peerError.Fatal || r.peerManager.HasMaxPeerCapacity() + r.logger.Error("peer error", + "peer", peerError.NodeID, + "err", peerError.Err, + "evicting", shouldEvict, + ) + if shouldEvict { + r.peerManager.Errored(peerError.NodeID, peerError.Err) + } else { + r.peerManager.processPeerEvent(ctx, PeerUpdate{ + NodeID: peerError.NodeID, + Status: PeerStatusBad, + }) + } - r.peerManager.Errored(peerError.NodeID, peerError.Err) case <-ctx.Done(): return }