mirror of
https://github.com/netbirdio/netbird.git
synced 2026-08-24 16:41:30 +02:00
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.
143 lines
3.5 KiB
Go
143 lines
3.5 KiB
Go
package peer
|
|
|
|
import (
|
|
"context"
|
|
"errors"
|
|
"net"
|
|
"net/netip"
|
|
"sync"
|
|
"sync/atomic"
|
|
|
|
log "github.com/sirupsen/logrus"
|
|
|
|
relayClient "github.com/netbirdio/netbird/shared/relay/client"
|
|
)
|
|
|
|
type RelayConnInfo struct {
|
|
relayedConn net.Conn
|
|
rosenpassPubKey []byte
|
|
rosenpassAddr string
|
|
}
|
|
|
|
type WorkerRelay struct {
|
|
peerCtx context.Context
|
|
log *log.Entry
|
|
isController bool
|
|
config ConnConfig
|
|
conn *Conn
|
|
relayManager *relayClient.Manager
|
|
|
|
relayedConn net.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()
|
|
|
|
err = w.relayManager.AddCloseListener(srv, w.onRelayClientDisconnected)
|
|
if err != nil {
|
|
log.Errorf("failed to add close listener: %s", err)
|
|
_ = relayedConn.Close()
|
|
return
|
|
}
|
|
|
|
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()
|
|
defer w.relayLock.Unlock()
|
|
if w.relayedConn == nil {
|
|
return
|
|
}
|
|
|
|
if err := w.relayedConn.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) onRelayClientDisconnected() {
|
|
go w.conn.onRelayDisconnected()
|
|
}
|