Files
netbird/client/internal/peer/worker_relay.go
Riccardo Manfrin f175e402c7 [client] stop offering to every peer when the relay transport drops (#7092)
* [client] stop offering to every peer when the relay transport drops

The relay transport is shared: one connection per relay server carries the
streams of every peer using it. When it drops, each of those peers gets a
Disconnected verdict from evalConnStatus even when ICE is still carrying its
traffic, because peerUsesRelay comes from HasRelayAddress(), which only reports
that management offered relay servers, not that we are connected to one. The
guard answers Disconnected with the aggressive retry, so every peer starts
sending offers over signal for a transport that no offer can restore: the relay
client's own guard is what reconnects it.

Feed relayManager.Ready() into the status inputs and return PartiallyConnected
when ICE is up and the missing side is the shared transport. That is the
existing "one path works, the other does not" branch, which retries three times
and then hourly instead of walking the exponential ladder forever.

Peers are not left waiting for the hourly tick: when the transport comes back,
Manager.onServerConnected notifies srWatcher, the guard resets the ticker to
800ms and iceState.reset() clears the hourly mode.

The verdict is unchanged when the transport is up but this peer is unreachable
over relay - it may have moved to another server, and only an offer carries its
new relay address - and in force-relay mode, where relay is the only transport.

* Renaming according to actual meanings

* Don't consider an in progress ICE as "partially connected"

when the relay is not..

* Aligns tests

* Address wrong comments
2026-10-05 15:59:31 +02:00

142 lines
3.4 KiB
Go

package peer
import (
"context"
"errors"
"net/netip"
"sync"
"sync/atomic"
log "github.com/sirupsen/logrus"
relayClient "github.com/netbirdio/netbird/shared/relay/client"
)
type RelayConnInfo struct {
relayedConn *relayClient.Conn
rosenpassPubKey []byte
rosenpassAddr string
}
type WorkerRelay struct {
peerCtx context.Context
log *log.Entry
isController bool
config ConnConfig
conn *Conn
relayManager *relayClient.Manager
relayedConn *relayClient.Conn
relayLock sync.Mutex
relaySupportedOnRemotePeer atomic.Bool
}
func NewWorkerRelay(ctx context.Context, log *log.Entry, ctrl bool, config ConnConfig, conn *Conn, relayManager *relayClient.Manager) *WorkerRelay {
r := &WorkerRelay{
peerCtx: ctx,
log: log,
isController: ctrl,
config: config,
conn: conn,
relayManager: relayManager,
}
return r
}
func (w *WorkerRelay) OnNewOffer(remoteOfferAnswer *OfferAnswer) {
if !w.isRelaySupported(remoteOfferAnswer) {
w.log.Infof("Relay is not supported by remote peer")
w.relaySupportedOnRemotePeer.Store(false)
return
}
w.relaySupportedOnRemotePeer.Store(true)
// the relayManager will return with error in case if the connection has lost with relay server
currentRelayAddress, _, err := w.relayManager.RelayInstanceAddress()
if err != nil {
w.log.Errorf("failed to handle new offer: %s", err)
return
}
srv := w.preferredRelayServer(currentRelayAddress, remoteOfferAnswer.RelaySrvAddress)
var serverIP netip.Addr
if srv == remoteOfferAnswer.RelaySrvAddress {
serverIP = remoteOfferAnswer.RelaySrvIP
}
relayedConn, err := w.relayManager.OpenConn(w.peerCtx, srv, w.config.Key, serverIP)
if err != nil {
if errors.Is(err, relayClient.ErrConnAlreadyExists) {
w.log.Debugf("handled offer by reusing existing relay connection")
return
}
w.log.Errorf("failed to open connection via Relay: %s", err)
return
}
w.relayLock.Lock()
w.relayedConn = relayedConn
w.relayLock.Unlock()
go w.watchRelayedConn(relayedConn)
w.log.Debugf("peer conn opened via Relay: %s", srv)
go w.conn.onRelayConnectionIsReady(RelayConnInfo{
relayedConn: relayedConn,
rosenpassPubKey: remoteOfferAnswer.RosenpassPubKey,
rosenpassAddr: remoteOfferAnswer.RosenpassAddr,
})
}
func (w *WorkerRelay) RelayInstanceAddress() (string, netip.Addr, error) {
return w.relayManager.RelayInstanceAddress()
}
func (w *WorkerRelay) IsRelayConnectionSupportedWithPeer() bool {
return w.relaySupportedOnRemotePeer.Load() && w.RelayIsSupportedLocally()
}
func (w *WorkerRelay) RelayIsSupportedLocally() bool {
return w.relayManager.HasRelayAddress()
}
func (w *WorkerRelay) IsTransportConnected() bool {
return w.relayManager.Ready()
}
func (w *WorkerRelay) CloseConn() {
w.relayLock.Lock()
conn := w.relayedConn
w.relayedConn = nil
w.relayLock.Unlock()
if conn == nil {
return
}
if err := conn.Close(); err != nil {
w.log.Warnf("failed to close relay connection: %v", err)
}
}
func (w *WorkerRelay) isRelaySupported(answer *OfferAnswer) bool {
if !w.relayManager.HasRelayAddress() {
return false
}
return answer.RelaySrvAddress != ""
}
func (w *WorkerRelay) preferredRelayServer(myRelayAddress, remoteRelayAddress string) string {
if w.isController {
return myRelayAddress
}
return remoteRelayAddress
}
func (w *WorkerRelay) watchRelayedConn(relayedConn *relayClient.Conn) {
<-relayedConn.Context().Done()
w.conn.onRelayDisconnected(relayedConn)
}