mirror of
https://github.com/netbirdio/netbird.git
synced 2026-10-07 14:09:07 +02:00
* [client] Coalesce peer list change notifications to the mobile listener Every peer state change spawned a goroutine to call the platform listener. During a reconnect storm this pinned hundreds of OS threads in JNI and let the UI call back into the engine from each of them. Deliver peer list changes from a single goroutine per listener and collapse pending changes into the latest count. * [client] Cover a pending wake-up when the peer list deliverer is replaced The replacement test waited for the old callback to finish before swapping listeners, so it never exercised the stop check that runs after a wake-up. Block the old callback, queue a peer list change and swap while it is blocked, then assert the old listener never sees the new count. * [client] Signal peer list deliverer exit and wait for it in the test The replacement test sampled the old listener after a fixed sleep, so a late stale delivery could slip past it. Close a done channel when the deliverer goroutine returns and let the test wait on it instead. * [client] Drop the test-only peer list deliverer exit channel The done channel was only read by the replacement test. Production code cannot wait on it, since joining the deliverer would block on a mobile callback. The tests now drive the deliverer loop directly and check that setListener and removeListener close its stop channel.
269 lines
6.0 KiB
Go
269 lines
6.0 KiB
Go
package peer
|
|
|
|
import (
|
|
"sync"
|
|
)
|
|
|
|
type notifier struct {
|
|
// publishLock orders state publication: it is held across computing the
|
|
// effective state and handing it to the listener, so a transition cannot
|
|
// overtake a newer one and leave the listener on a stale state.
|
|
publishLock sync.Mutex
|
|
serverStateLock sync.Mutex
|
|
listenersLock sync.Mutex
|
|
listener Listener
|
|
peerListWake chan struct{}
|
|
peerListStop chan struct{}
|
|
currentClientState bool
|
|
lastNotification ClientState
|
|
lastNumberOfPeers int
|
|
lastFqdnAddress string
|
|
lastIPAddress string
|
|
networkAvailable bool
|
|
}
|
|
|
|
func newNotifier() *notifier {
|
|
return ¬ifier{
|
|
networkAvailable: true,
|
|
}
|
|
}
|
|
|
|
// effectiveState maps the computed state to what listeners should see:
|
|
// while the OS reports no usable network, "Connecting" would be a lie —
|
|
// connection attempts are suspended — so it is reported as NoNetwork.
|
|
// Caller must hold serverStateLock.
|
|
func (n *notifier) effectiveState(state ClientState) ClientState {
|
|
if !n.networkAvailable && state == ClientStateConnecting {
|
|
return ClientStateNoNetwork
|
|
}
|
|
return state
|
|
}
|
|
|
|
// setNetworkAvailable records the OS network availability and re-notifies
|
|
// the listener when the flag flips the effective state (Connecting <->
|
|
// NoNetwork).
|
|
func (n *notifier) setNetworkAvailable(available bool) {
|
|
n.publishLock.Lock()
|
|
defer n.publishLock.Unlock()
|
|
|
|
n.serverStateLock.Lock()
|
|
if n.networkAvailable == available {
|
|
n.serverStateLock.Unlock()
|
|
return
|
|
}
|
|
previous := n.effectiveState(n.lastNotification)
|
|
n.networkAvailable = available
|
|
current := n.effectiveState(n.lastNotification)
|
|
n.serverStateLock.Unlock()
|
|
|
|
if previous != current {
|
|
n.notify(current)
|
|
}
|
|
}
|
|
|
|
func (n *notifier) setListener(listener Listener) {
|
|
n.serverStateLock.Lock()
|
|
lastNotification := n.effectiveState(n.lastNotification)
|
|
fqdnAddress := n.lastFqdnAddress
|
|
address := n.lastIPAddress
|
|
n.serverStateLock.Unlock()
|
|
|
|
n.listenersLock.Lock()
|
|
defer n.listenersLock.Unlock()
|
|
|
|
n.stopPeerListDelivererLocked()
|
|
n.listener = listener
|
|
|
|
listener.OnAddressChanged(fqdnAddress, address)
|
|
notifyListener(listener, lastNotification)
|
|
n.startPeerListDelivererLocked(listener)
|
|
n.wakePeerListDelivererLocked()
|
|
}
|
|
|
|
func (n *notifier) removeListener() {
|
|
n.listenersLock.Lock()
|
|
defer n.listenersLock.Unlock()
|
|
n.stopPeerListDelivererLocked()
|
|
n.listener = nil
|
|
}
|
|
|
|
func (n *notifier) updateServerStates(mgmState bool, signalState bool) {
|
|
n.publishLock.Lock()
|
|
defer n.publishLock.Unlock()
|
|
|
|
n.serverStateLock.Lock()
|
|
calculatedState := n.calculateState(mgmState, signalState)
|
|
|
|
if !n.isServerStateChanged(calculatedState) {
|
|
n.serverStateLock.Unlock()
|
|
return
|
|
}
|
|
|
|
n.lastNotification = calculatedState
|
|
effective := n.effectiveState(calculatedState)
|
|
n.serverStateLock.Unlock()
|
|
|
|
n.notify(effective)
|
|
}
|
|
|
|
func (n *notifier) clientStart() {
|
|
n.publishLock.Lock()
|
|
defer n.publishLock.Unlock()
|
|
|
|
n.serverStateLock.Lock()
|
|
n.currentClientState = true
|
|
n.lastNotification = ClientStateConnecting
|
|
effective := n.effectiveState(ClientStateConnecting)
|
|
n.serverStateLock.Unlock()
|
|
|
|
n.notify(effective)
|
|
}
|
|
|
|
func (n *notifier) clientStop() {
|
|
n.publishLock.Lock()
|
|
defer n.publishLock.Unlock()
|
|
|
|
n.serverStateLock.Lock()
|
|
n.currentClientState = false
|
|
n.lastNotification = ClientStateDisconnected
|
|
n.serverStateLock.Unlock()
|
|
|
|
n.notify(ClientStateDisconnected)
|
|
}
|
|
|
|
func (n *notifier) clientTearDown() {
|
|
n.publishLock.Lock()
|
|
defer n.publishLock.Unlock()
|
|
|
|
n.serverStateLock.Lock()
|
|
n.currentClientState = false
|
|
n.lastNotification = ClientStateDisconnecting
|
|
n.serverStateLock.Unlock()
|
|
|
|
n.notify(ClientStateDisconnecting)
|
|
}
|
|
|
|
func (n *notifier) isServerStateChanged(newState ClientState) bool {
|
|
return n.lastNotification != newState
|
|
}
|
|
|
|
func (n *notifier) notify(state ClientState) {
|
|
n.listenersLock.Lock()
|
|
listener := n.listener
|
|
n.listenersLock.Unlock()
|
|
|
|
if listener == nil {
|
|
return
|
|
}
|
|
|
|
notifyListener(listener, state)
|
|
}
|
|
|
|
func (n *notifier) calculateState(managementConn, signalConn bool) ClientState {
|
|
if managementConn && signalConn {
|
|
return ClientStateConnected
|
|
}
|
|
|
|
if !managementConn && !signalConn && !n.currentClientState {
|
|
return ClientStateDisconnected
|
|
}
|
|
|
|
if n.lastNotification == ClientStateDisconnecting {
|
|
return ClientStateDisconnecting
|
|
}
|
|
|
|
return ClientStateConnecting
|
|
}
|
|
|
|
func (n *notifier) peerListChanged(numOfPeers int) {
|
|
n.serverStateLock.Lock()
|
|
n.lastNumberOfPeers = numOfPeers
|
|
n.serverStateLock.Unlock()
|
|
|
|
n.listenersLock.Lock()
|
|
defer n.listenersLock.Unlock()
|
|
n.wakePeerListDelivererLocked()
|
|
}
|
|
|
|
func (n *notifier) startPeerListDelivererLocked(listener Listener) {
|
|
wake := make(chan struct{}, 1)
|
|
stop := make(chan struct{})
|
|
n.peerListWake = wake
|
|
n.peerListStop = stop
|
|
go n.deliverPeerListChanges(listener, wake, stop)
|
|
}
|
|
|
|
func (n *notifier) stopPeerListDelivererLocked() {
|
|
if n.peerListStop == nil {
|
|
return
|
|
}
|
|
close(n.peerListStop)
|
|
n.peerListStop = nil
|
|
n.peerListWake = nil
|
|
}
|
|
|
|
func (n *notifier) wakePeerListDelivererLocked() {
|
|
if n.peerListWake == nil {
|
|
return
|
|
}
|
|
select {
|
|
case n.peerListWake <- struct{}{}:
|
|
default:
|
|
}
|
|
}
|
|
|
|
func (n *notifier) deliverPeerListChanges(listener Listener, wake <-chan struct{}, stop <-chan struct{}) {
|
|
for {
|
|
select {
|
|
case <-stop:
|
|
return
|
|
case <-wake:
|
|
}
|
|
select {
|
|
case <-stop:
|
|
return
|
|
default:
|
|
}
|
|
|
|
n.serverStateLock.Lock()
|
|
numOfPeers := n.lastNumberOfPeers
|
|
n.serverStateLock.Unlock()
|
|
|
|
listener.OnPeersListChanged(numOfPeers)
|
|
}
|
|
}
|
|
|
|
func (n *notifier) localAddressChanged(fqdn, address string) {
|
|
n.serverStateLock.Lock()
|
|
n.lastFqdnAddress = fqdn
|
|
n.lastIPAddress = address
|
|
n.serverStateLock.Unlock()
|
|
|
|
n.listenersLock.Lock()
|
|
listener := n.listener
|
|
n.listenersLock.Unlock()
|
|
|
|
if listener == nil {
|
|
return
|
|
}
|
|
|
|
listener.OnAddressChanged(fqdn, address)
|
|
}
|
|
|
|
func notifyListener(l Listener, state ClientState) {
|
|
// legacy per-state callbacks; NoNetwork is delivered only via
|
|
// OnStateChanged below
|
|
switch state {
|
|
case ClientStateDisconnected:
|
|
l.OnDisconnected()
|
|
case ClientStateConnected:
|
|
l.OnConnected()
|
|
case ClientStateConnecting:
|
|
l.OnConnecting()
|
|
case ClientStateDisconnecting:
|
|
l.OnDisconnecting()
|
|
}
|
|
|
|
l.OnStateChanged(state)
|
|
}
|