mirror of
https://github.com/netbirdio/netbird.git
synced 2026-10-06 13:39:07 +02:00
* [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
215 lines
6.8 KiB
Go
215 lines
6.8 KiB
Go
package guard
|
|
|
|
import (
|
|
"context"
|
|
"time"
|
|
|
|
"github.com/cenkalti/backoff/v4"
|
|
log "github.com/sirupsen/logrus"
|
|
)
|
|
|
|
// ConnStatus represents the connection state as seen by the guard.
|
|
type ConnStatus int
|
|
|
|
const (
|
|
// ConnStatusDisconnected means neither ICE nor Relay is connected.
|
|
ConnStatusDisconnected ConnStatus = iota
|
|
// ConnStatusPartiallyConnected means one transport is usable and the other is not:
|
|
// relay connected with ICE down, or ICE connected with the shared relay transport down.
|
|
ConnStatusPartiallyConnected
|
|
// ConnStatusConnected means all required connections are established.
|
|
ConnStatusConnected
|
|
)
|
|
|
|
type connStatusFunc func() ConnStatus
|
|
|
|
// NetworkWatcher is the availability view the guard gates reconnects on.
|
|
type NetworkWatcher interface {
|
|
IsOnline() bool
|
|
Changed() <-chan struct{}
|
|
}
|
|
|
|
// Guard is responsible for the reconnection logic.
|
|
// It will trigger to send an offer to the peer then has connection issues.
|
|
// Watch these events:
|
|
// - Relay client reconnected to home server
|
|
// - Signal server connection state changed
|
|
// - ICE connection disconnected
|
|
// - Relayed connection disconnected
|
|
// - ICE candidate changes
|
|
type Guard struct {
|
|
log *log.Entry
|
|
isConnectedOnAllWay connStatusFunc
|
|
timeout time.Duration
|
|
srWatcher *SRWatcher
|
|
// netWatcher gates reconnect attempts on OS-reported network availability;
|
|
// nil disables gating.
|
|
netWatcher NetworkWatcher
|
|
relayedConnDisconnected chan struct{}
|
|
iCEConnDisconnected chan struct{}
|
|
}
|
|
|
|
// NewGuard creates a reconnection guard for a peer connection. A nil netWatcher
|
|
// disables network availability gating.
|
|
func NewGuard(log *log.Entry, isConnectedFn connStatusFunc, timeout time.Duration, srWatcher *SRWatcher, netWatcher NetworkWatcher) *Guard {
|
|
return &Guard{
|
|
log: log,
|
|
isConnectedOnAllWay: isConnectedFn,
|
|
timeout: timeout,
|
|
srWatcher: srWatcher,
|
|
netWatcher: netWatcher,
|
|
relayedConnDisconnected: make(chan struct{}, 1),
|
|
iCEConnDisconnected: make(chan struct{}, 1),
|
|
}
|
|
}
|
|
|
|
func (g *Guard) Start(ctx context.Context, eventCallback func()) {
|
|
g.log.Infof("starting guard for reconnection with MaxInterval: %s", g.timeout)
|
|
g.reconnectLoopWithRetry(ctx, eventCallback)
|
|
}
|
|
|
|
func (g *Guard) SetRelayedConnDisconnected() {
|
|
select {
|
|
case g.relayedConnDisconnected <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
func (g *Guard) SetICEConnDisconnected() {
|
|
select {
|
|
case g.iCEConnDisconnected <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
// reconnectLoopWithRetry periodically checks the connection status and sends offers to re-establish connectivity.
|
|
//
|
|
// Behavior depends on the connection state reported by isConnectedOnAllWay:
|
|
// - Connected: no action, the peer is fully reachable.
|
|
// - Disconnected (neither ICE nor Relay): retries aggressively with exponential backoff (800ms doubling
|
|
// up to timeout), never gives up. This ensures rapid recovery when the peer has no connectivity at all.
|
|
// - PartiallyConnected (one transport usable, the other not): retries up to 3 times
|
|
// with exponential backoff, then switches to one attempt per hour. This limits
|
|
// signaling traffic while the peer still has a working path.
|
|
//
|
|
// External events (relay/ICE disconnect, signal/relay reconnect, candidate changes) reset the retry
|
|
// counter and backoff ticker, giving ICE a fresh chance after network conditions change.
|
|
func (g *Guard) reconnectLoopWithRetry(ctx context.Context, callback func()) {
|
|
srReconnectedChan := g.srWatcher.NewListener()
|
|
defer g.srWatcher.RemoveListener(srReconnectedChan)
|
|
|
|
ticker := g.initialTicker(ctx)
|
|
defer func() {
|
|
// If backoff.Ticker.send is blocked, context.Done will not close the Ticker goroutine.
|
|
// We have to explicitly call Stop, even if we use backoff.WithContext.
|
|
ticker.Stop()
|
|
}()
|
|
|
|
tickerChannel := ticker.C
|
|
|
|
iceState := &iceRetryState{log: g.log}
|
|
defer iceState.reset()
|
|
|
|
var netChanged <-chan struct{}
|
|
if g.netWatcher != nil {
|
|
netChanged = g.netWatcher.Changed()
|
|
}
|
|
|
|
for {
|
|
select {
|
|
case <-tickerChannel:
|
|
// skip attempts while the OS reports no usable network; the
|
|
// netChanged case below resumes the loop once it returns
|
|
if g.netWatcher != nil && !g.netWatcher.IsOnline() {
|
|
continue
|
|
}
|
|
switch g.isConnectedOnAllWay() {
|
|
case ConnStatusConnected:
|
|
// all good, nothing to do
|
|
case ConnStatusDisconnected:
|
|
callback()
|
|
case ConnStatusPartiallyConnected:
|
|
if iceState.shouldRetry() {
|
|
callback()
|
|
} else {
|
|
iceState.enterHourlyMode()
|
|
ticker.Stop()
|
|
tickerChannel = iceState.hourlyC()
|
|
}
|
|
}
|
|
|
|
case <-g.relayedConnDisconnected:
|
|
g.log.Debugf("Relay connection changed, reset reconnection ticker")
|
|
ticker.Stop()
|
|
ticker = g.newReconnectTicker(ctx)
|
|
tickerChannel = ticker.C
|
|
iceState.reset()
|
|
|
|
case <-g.iCEConnDisconnected:
|
|
g.log.Debugf("ICE connection changed, reset reconnection ticker")
|
|
ticker.Stop()
|
|
ticker = g.newReconnectTicker(ctx)
|
|
tickerChannel = ticker.C
|
|
iceState.reset()
|
|
|
|
case <-srReconnectedChan:
|
|
g.log.Debugf("has network changes, reset reconnection ticker")
|
|
ticker.Stop()
|
|
ticker = g.newReconnectTicker(ctx)
|
|
tickerChannel = ticker.C
|
|
iceState.reset()
|
|
|
|
case <-netChanged:
|
|
// Re-arm for the next transition before acting on this one.
|
|
netChanged = g.netWatcher.Changed()
|
|
if !g.netWatcher.IsOnline() {
|
|
continue
|
|
}
|
|
// Ticks skipped while offline drove the backoff towards its
|
|
// maximum without ever attempting, and left the ICE budget
|
|
// frozen — possibly in hourly mode. Recover on our own so the
|
|
// peer does not depend on a signal or relay event that never
|
|
// comes when both stayed up across the outage.
|
|
g.log.Debugf("network is back, reset reconnection ticker")
|
|
ticker.Stop()
|
|
ticker = g.newReconnectTicker(ctx)
|
|
tickerChannel = ticker.C
|
|
iceState.reset()
|
|
|
|
case <-ctx.Done():
|
|
g.log.Debugf("context is done, stop reconnect loop")
|
|
return
|
|
}
|
|
}
|
|
}
|
|
|
|
// initialTicker give chance to the peer to establish the initial connection.
|
|
func (g *Guard) initialTicker(ctx context.Context) *backoff.Ticker {
|
|
bo := backoff.WithContext(&backoff.ExponentialBackOff{
|
|
InitialInterval: 3 * time.Second,
|
|
RandomizationFactor: 0.1,
|
|
Multiplier: 2,
|
|
MaxInterval: g.timeout,
|
|
Stop: backoff.Stop,
|
|
Clock: backoff.SystemClock,
|
|
}, ctx)
|
|
|
|
return backoff.NewTicker(bo)
|
|
}
|
|
|
|
func (g *Guard) newReconnectTicker(ctx context.Context) *backoff.Ticker {
|
|
bo := backoff.WithContext(&backoff.ExponentialBackOff{
|
|
InitialInterval: 800 * time.Millisecond,
|
|
RandomizationFactor: 0.1,
|
|
Multiplier: 2,
|
|
MaxInterval: g.timeout,
|
|
Stop: backoff.Stop,
|
|
Clock: backoff.SystemClock,
|
|
}, ctx)
|
|
|
|
ticker := backoff.NewTicker(bo)
|
|
<-ticker.C // consume the initial tick what is happening right after the ticker has been created
|
|
|
|
return ticker
|
|
}
|