[client] Extract peer status recorder into its own package

Move the Status recorder and its state types out of the peer package
into client/internal/peer/status, split by struct across recorder.go,
peer_state.go, full_status.go, events.go, notifier.go and route.go
instead of one 1600-line file. Rename the type Status -> Recorder
(NewRecorder already implied it; avoids status.Status stutter). Split
conn_status.go: the ConnStatus type and its constants move to the status
package, connStatusInputs stays with the peer event loop.

The peer package references the status package directly; a transitional
status_alias.go re-exports the moved symbols for the ~50 external callers
still using peer.Status/State/ConnStatus, to be removed once they are
migrated.
This commit is contained in:
Zoltan Papp
2026-08-05 16:14:11 +02:00
committed by Zoltán Papp
parent 71519f1b5d
commit 9d92bc8abc
18 changed files with 429 additions and 362 deletions
+16 -15
View File
@@ -22,6 +22,7 @@ import (
"github.com/netbirdio/netbird/client/internal/peer/guard" "github.com/netbirdio/netbird/client/internal/peer/guard"
icemaker "github.com/netbirdio/netbird/client/internal/peer/ice" icemaker "github.com/netbirdio/netbird/client/internal/peer/ice"
"github.com/netbirdio/netbird/client/internal/peer/id" "github.com/netbirdio/netbird/client/internal/peer/id"
"github.com/netbirdio/netbird/client/internal/peer/status"
"github.com/netbirdio/netbird/client/internal/peer/worker" "github.com/netbirdio/netbird/client/internal/peer/worker"
"github.com/netbirdio/netbird/client/internal/portforward" "github.com/netbirdio/netbird/client/internal/portforward"
"github.com/netbirdio/netbird/client/internal/rosenpass" "github.com/netbirdio/netbird/client/internal/rosenpass"
@@ -47,7 +48,7 @@ type MetricsRecorder interface {
} }
type ServiceDependencies struct { type ServiceDependencies struct {
StatusRecorder *Status StatusRecorder *status.Recorder
Signaler *Signaler Signaler *Signaler
IFaceDiscover stdnet.ExternalIFaceDiscover IFaceDiscover stdnet.ExternalIFaceDiscover
RelayManager *relayClient.Manager RelayManager *relayClient.Manager
@@ -106,7 +107,7 @@ type Conn struct {
ctx context.Context ctx context.Context
ctxCancel context.CancelFunc ctxCancel context.CancelFunc
config ConnConfig config ConnConfig
statusRecorder *Status statusRecorder *status.Recorder
signaler *Signaler signaler *Signaler
iFaceDiscover stdnet.ExternalIFaceDiscover iFaceDiscover stdnet.ExternalIFaceDiscover
relayManager *relayClient.Manager relayManager *relayClient.Manager
@@ -244,10 +245,10 @@ func (conn *Conn) open(engineCtx context.Context, firstPacket []byte) error {
conn.pendingFirstPacket = slices.Clone(firstPacket) conn.pendingFirstPacket = slices.Clone(firstPacket)
} }
peerState := State{ peerState := status.State{
PubKey: conn.config.Key, PubKey: conn.config.Key,
ConnStatusUpdate: time.Now(), ConnStatusUpdate: time.Now(),
ConnStatus: StatusConnecting, ConnStatus: status.StatusConnecting,
Mux: new(sync.RWMutex), Mux: new(sync.RWMutex),
} }
if err := conn.statusRecorder.UpdatePeerState(peerState); err != nil { if err := conn.statusRecorder.UpdatePeerState(peerState); err != nil {
@@ -346,7 +347,7 @@ func (conn *Conn) WgConfig() WgConfig {
// IsConnected returns true if the peer is connected // IsConnected returns true if the peer is connected
func (conn *Conn) IsConnected() bool { func (conn *Conn) IsConnected() bool {
return conn.evalStatus() == StatusConnected return conn.evalStatus() == status.StatusConnected
} }
func (conn *Conn) GetKey() string { func (conn *Conn) GetKey() string {
@@ -473,7 +474,7 @@ func (conn *Conn) teardown(mb *mailbox, leftover []event, signalToRemote bool, d
conn.Log.Errorf("failed to remove wg endpoint: %v", err) conn.Log.Errorf("failed to remove wg endpoint: %v", err)
} }
if conn.evalStatus() == StatusConnected && conn.onDisconnected != nil { if conn.evalStatus() == status.StatusConnected && conn.onDisconnected != nil {
conn.onDisconnected(conn.config.WgConfig.RemoteKey) conn.onDisconnected(conn.config.WgConfig.RemoteKey)
} }
@@ -716,7 +717,7 @@ func (conn *Conn) handleICEDisconnected(sessionChanged bool) {
conn.metricsStages.Disconnected() conn.metricsStages.Disconnected()
} }
peerState := State{ peerState := status.State{
PubKey: conn.config.Key, PubKey: conn.config.Key,
ConnStatus: conn.evalStatus(), ConnStatus: conn.evalStatus(),
Relayed: conn.isRelayed(), Relayed: conn.isRelayed(),
@@ -824,7 +825,7 @@ func (conn *Conn) handleRelayDisconnected() {
conn.metricsStages.Disconnected() conn.metricsStages.Disconnected()
} }
peerState := State{ peerState := status.State{
PubKey: conn.config.Key, PubKey: conn.config.Key,
ConnStatus: conn.evalStatus(), ConnStatus: conn.evalStatus(),
Relayed: conn.isRelayed(), Relayed: conn.isRelayed(),
@@ -957,7 +958,7 @@ func (conn *Conn) injectPendingFirstPacket(proxy wgproxy.Proxy, directConn net.C
} }
func (conn *Conn) updateRelayStatus(relayServerAddr string, rosenpassPubKey []byte, updateTime time.Time) { func (conn *Conn) updateRelayStatus(relayServerAddr string, rosenpassPubKey []byte, updateTime time.Time) {
peerState := State{ peerState := status.State{
PubKey: conn.config.Key, PubKey: conn.config.Key,
ConnStatusUpdate: updateTime, ConnStatusUpdate: updateTime,
ConnStatus: conn.evalStatus(), ConnStatus: conn.evalStatus(),
@@ -973,7 +974,7 @@ func (conn *Conn) updateRelayStatus(relayServerAddr string, rosenpassPubKey []by
} }
func (conn *Conn) updateIceState(iceConnInfo ICEConnInfo, updateTime time.Time) { func (conn *Conn) updateIceState(iceConnInfo ICEConnInfo, updateTime time.Time) {
peerState := State{ peerState := status.State{
PubKey: conn.config.Key, PubKey: conn.config.Key,
ConnStatusUpdate: updateTime, ConnStatusUpdate: updateTime,
ConnStatus: conn.evalStatus(), ConnStatus: conn.evalStatus(),
@@ -996,9 +997,9 @@ func (conn *Conn) setStatusToDisconnected() {
conn.statusICE.SetDisconnected() conn.statusICE.SetDisconnected()
conn.currentConnPriority = conntype.None conn.currentConnPriority = conntype.None
peerState := State{ peerState := status.State{
PubKey: conn.config.Key, PubKey: conn.config.Key,
ConnStatus: StatusIdle, ConnStatus: status.StatusIdle,
ConnStatusUpdate: time.Now(), ConnStatusUpdate: time.Now(),
Mux: new(sync.RWMutex), Mux: new(sync.RWMutex),
} }
@@ -1034,12 +1035,12 @@ func (conn *Conn) isRelayed() bool {
} }
} }
func (conn *Conn) evalStatus() ConnStatus { func (conn *Conn) evalStatus() status.ConnStatus {
if conn.statusRelay.Get() == worker.StatusConnected || conn.statusICE.Get() == worker.StatusConnected { if conn.statusRelay.Get() == worker.StatusConnected || conn.statusICE.Get() == worker.StatusConnected {
return StatusConnected return status.StatusConnected
} }
return StatusConnecting return status.StatusConnecting
} }
// isConnectedOnAllWay evaluates the overall connection status based on ICE and Relay transports. // isConnectedOnAllWay evaluates the overall connection status based on ICE and Relay transports.
-30
View File
@@ -1,18 +1,5 @@
package peer package peer
import (
log "github.com/sirupsen/logrus"
)
const (
// StatusIdle indicate the peer is in disconnected state
StatusIdle ConnStatus = iota
// StatusConnecting indicate the peer is in connecting state
StatusConnecting
// StatusConnected indicate the peer is in connected state
StatusConnected
)
// connStatusInputs is the primitive-valued snapshot of the state that drives the // connStatusInputs is the primitive-valued snapshot of the state that drives the
// tri-state connection classification. Extracted so the decision logic can be unit-tested // tri-state connection classification. Extracted so the decision logic can be unit-tested
// without constructing full Worker/Handshaker objects. // without constructing full Worker/Handshaker objects.
@@ -25,20 +12,3 @@ type connStatusInputs struct {
iceStatusConnecting bool // statusICE is anything other than Disconnected iceStatusConnecting bool // statusICE is anything other than Disconnected
iceInProgress bool // a negotiation is currently in flight iceInProgress bool // a negotiation is currently in flight
} }
// ConnStatus describe the status of a peer's connection
type ConnStatus int32
func (s ConnStatus) String() string {
switch s {
case StatusConnecting:
return "Connecting"
case StatusConnected:
return "Connected"
case StatusIdle:
return "Idle"
default:
log.Errorf("unknown status: %d", s)
return "INVALID_PEER_CONNECTION_STATUS"
}
}
+2 -1
View File
@@ -15,6 +15,7 @@ import (
"github.com/netbirdio/netbird/client/iface" "github.com/netbirdio/netbird/client/iface"
"github.com/netbirdio/netbird/client/internal/peer/guard" "github.com/netbirdio/netbird/client/internal/peer/guard"
"github.com/netbirdio/netbird/client/internal/peer/ice" "github.com/netbirdio/netbird/client/internal/peer/ice"
"github.com/netbirdio/netbird/client/internal/peer/status"
"github.com/netbirdio/netbird/client/internal/stdnet" "github.com/netbirdio/netbird/client/internal/stdnet"
"github.com/netbirdio/netbird/util" "github.com/netbirdio/netbird/util"
) )
@@ -69,7 +70,7 @@ func TestConn_GetKey(t *testing.T) {
func TestConn_DiscardMessagesWhenNotOpened(t *testing.T) { func TestConn_DiscardMessagesWhenNotOpened(t *testing.T) {
swWatcher := guard.NewSRWatcher(nil, nil, nil, connConf.ICEConfig) swWatcher := guard.NewSRWatcher(nil, nil, nil, connConf.ICEConfig)
sd := ServiceDependencies{ sd := ServiceDependencies{
StatusRecorder: NewRecorder("https://mgm"), StatusRecorder: status.NewRecorder("https://mgm"),
SrWatcher: swWatcher, SrWatcher: swWatcher,
} }
conn, err := NewConn(connConf, sd) conn, err := NewConn(connConf, sd)
-11
View File
@@ -1,11 +0,0 @@
package peer
// Listener is a callback type about the NetBird network connection state
type Listener interface {
OnConnected()
OnDisconnected()
OnConnecting()
OnDisconnecting()
OnAddressChanged(string, string)
OnPeersListChanged(int)
}
+4 -2
View File
@@ -6,11 +6,13 @@ import (
"time" "time"
log "github.com/sirupsen/logrus" log "github.com/sirupsen/logrus"
"github.com/netbirdio/netbird/client/internal/peer/status"
) )
type stateDump struct { type stateDump struct {
log *log.Entry log *log.Entry
status *Status status *status.Recorder
key string key string
sentOffer int sentOffer int
@@ -26,7 +28,7 @@ type stateDump struct {
mu sync.Mutex mu sync.Mutex
} }
func newStateDump(key string, log *log.Entry, statusRecorder *Status) *stateDump { func newStateDump(key string, log *log.Entry, statusRecorder *status.Recorder) *stateDump {
return &stateDump{ return &stateDump{
log: log, log: log,
status: statusRecorder, status: statusRecorder,
@@ -0,0 +1,31 @@
package status
import (
log "github.com/sirupsen/logrus"
)
const (
// StatusIdle indicate the peer is in disconnected state
StatusIdle ConnStatus = iota
// StatusConnecting indicate the peer is in connecting state
StatusConnecting
// StatusConnected indicate the peer is in connected state
StatusConnected
)
// ConnStatus describe the status of a peer's connection
type ConnStatus int32
func (s ConnStatus) String() string {
switch s {
case StatusConnecting:
return "Connecting"
case StatusConnected:
return "Connected"
case StatusIdle:
return "Idle"
default:
log.Errorf("unknown status: %d", s)
return "INVALID_PEER_CONNECTION_STATUS"
}
}
@@ -1,4 +1,4 @@
package peer package status
import ( import (
"testing" "testing"
+48
View File
@@ -0,0 +1,48 @@
package status
import (
"slices"
"sync"
"github.com/netbirdio/netbird/client/proto"
)
type EventQueue struct {
maxSize int
events []*proto.SystemEvent
mutex sync.RWMutex
}
func NewEventQueue(size int) *EventQueue {
return &EventQueue{
maxSize: size,
events: make([]*proto.SystemEvent, 0, size),
}
}
func (q *EventQueue) Add(event *proto.SystemEvent) {
q.mutex.Lock()
defer q.mutex.Unlock()
q.events = append(q.events, event)
if len(q.events) > q.maxSize {
q.events = q.events[len(q.events)-q.maxSize:]
}
}
func (q *EventQueue) GetAll() []*proto.SystemEvent {
q.mutex.RLock()
defer q.mutex.RUnlock()
return slices.Clone(q.events)
}
type EventSubscription struct {
id string
events chan *proto.SystemEvent
}
func (s *EventSubscription) Events() <-chan *proto.SystemEvent {
return s.events
}
+122
View File
@@ -0,0 +1,122 @@
package status
import (
"golang.org/x/exp/maps"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/timestamppb"
"github.com/netbirdio/netbird/client/internal/relay"
"github.com/netbirdio/netbird/client/proto"
)
// FullStatus contains the full state held by the Recorder instance
type FullStatus struct {
Peers []State
ManagementState ManagementState
SignalState SignalState
LocalPeerState LocalPeerState
RosenpassState RosenpassState
Relays []relay.ProbeResult
NSGroupStates []NSGroupState
NumOfForwardingRules int
LazyConnectionEnabled bool
Events []*proto.SystemEvent
}
// ToProto converts FullStatus to proto.FullStatus.
func (fs FullStatus) ToProto() *proto.FullStatus {
pbFullStatus := proto.FullStatus{
ManagementState: &proto.ManagementState{},
SignalState: &proto.SignalState{},
LocalPeerState: &proto.LocalPeerState{},
Peers: []*proto.PeerState{},
}
pbFullStatus.ManagementState.URL = fs.ManagementState.URL
pbFullStatus.ManagementState.Connected = fs.ManagementState.Connected
if err := fs.ManagementState.Error; err != nil {
pbFullStatus.ManagementState.Error = err.Error()
}
pbFullStatus.SignalState.URL = fs.SignalState.URL
pbFullStatus.SignalState.Connected = fs.SignalState.Connected
if err := fs.SignalState.Error; err != nil {
pbFullStatus.SignalState.Error = err.Error()
}
pbFullStatus.LocalPeerState.IP = fs.LocalPeerState.IP
pbFullStatus.LocalPeerState.Ipv6 = fs.LocalPeerState.IPv6
pbFullStatus.LocalPeerState.PubKey = fs.LocalPeerState.PubKey
pbFullStatus.LocalPeerState.KernelInterface = fs.LocalPeerState.KernelInterface
pbFullStatus.LocalPeerState.Fqdn = fs.LocalPeerState.FQDN
pbFullStatus.LocalPeerState.WgPort = int32(fs.LocalPeerState.WgPort)
pbFullStatus.LocalPeerState.RosenpassPermissive = fs.RosenpassState.Permissive
pbFullStatus.LocalPeerState.RosenpassEnabled = fs.RosenpassState.Enabled
pbFullStatus.NumberOfForwardingRules = int32(fs.NumOfForwardingRules)
pbFullStatus.LazyConnectionEnabled = fs.LazyConnectionEnabled
pbFullStatus.LocalPeerState.Networks = maps.Keys(fs.LocalPeerState.Routes)
for _, peerState := range fs.Peers {
networks := maps.Keys(peerState.GetRoutes())
pbPeerState := &proto.PeerState{
IP: peerState.IP,
Ipv6: peerState.IPv6,
PubKey: peerState.PubKey,
ConnStatus: peerState.ConnStatus.String(),
ConnStatusUpdate: timestamppb.New(peerState.ConnStatusUpdate),
Relayed: peerState.Relayed,
LocalIceCandidateType: peerState.LocalIceCandidateType,
RemoteIceCandidateType: peerState.RemoteIceCandidateType,
LocalIceCandidateEndpoint: peerState.LocalIceCandidateEndpoint,
RemoteIceCandidateEndpoint: peerState.RemoteIceCandidateEndpoint,
RelayAddress: peerState.RelayServerAddress,
Fqdn: peerState.FQDN,
LastWireguardHandshake: timestamppb.New(peerState.LastWireguardHandshake),
BytesRx: peerState.BytesRx,
BytesTx: peerState.BytesTx,
RosenpassEnabled: peerState.RosenpassEnabled,
Networks: networks,
Latency: durationpb.New(peerState.Latency),
SshHostKey: peerState.SSHHostKey,
}
pbFullStatus.Peers = append(pbFullStatus.Peers, pbPeerState)
}
for _, relayState := range fs.Relays {
pbRelayState := &proto.RelayState{
URI: relayState.URI,
Available: relayState.Err == nil,
Transport: relayState.Transport,
}
if err := relayState.Err; err != nil {
pbRelayState.Error = err.Error()
}
pbFullStatus.Relays = append(pbFullStatus.Relays, pbRelayState)
}
for _, dnsState := range fs.NSGroupStates {
var err string
if dnsState.Error != nil {
err = dnsState.Error.Error()
}
var servers []string
for _, server := range dnsState.Servers {
servers = append(servers, server.String())
}
pbDnsState := &proto.NSGroupState{
Servers: servers,
Domains: dnsState.Domains,
Enabled: dnsState.Enabled,
Error: err,
}
pbFullStatus.DnsServers = append(pbFullStatus.DnsServers, pbDnsState)
}
pbFullStatus.Events = fs.Events
return &pbFullStatus
}
@@ -1,4 +1,4 @@
package peer package status
import ( import (
"sync" "sync"
@@ -11,6 +11,16 @@ const (
stateDisconnecting stateDisconnecting
) )
// Listener is a callback type about the NetBird network connection state
type Listener interface {
OnConnected()
OnDisconnected()
OnConnecting()
OnDisconnecting()
OnAddressChanged(string, string)
OnPeersListChanged(int)
}
type notifier struct { type notifier struct {
serverStateLock sync.Mutex serverStateLock sync.Mutex
listenersLock sync.Mutex listenersLock sync.Mutex
@@ -1,4 +1,4 @@
package peer package status
import ( import (
"sync" "sync"
+63
View File
@@ -0,0 +1,63 @@
package status
import (
"sync"
"time"
"golang.org/x/exp/maps"
)
// State contains the latest state of a peer
type State struct {
Mux *sync.RWMutex
IP string
IPv6 string
PubKey string
FQDN string
ConnStatus ConnStatus
ConnStatusUpdate time.Time
Relayed bool
LocalIceCandidateType string
RemoteIceCandidateType string
LocalIceCandidateEndpoint string
RemoteIceCandidateEndpoint string
RelayServerAddress string
LastWireguardHandshake time.Time
BytesTx int64
BytesRx int64
Latency time.Duration
RosenpassEnabled bool
SSHHostKey []byte
routes map[string]struct{}
}
// AddRoute add a single route to routes map
func (s *State) AddRoute(network string) {
s.Mux.Lock()
defer s.Mux.Unlock()
if s.routes == nil {
s.routes = make(map[string]struct{})
}
s.routes[network] = struct{}{}
}
// SetRoutes set state routes
func (s *State) SetRoutes(routes map[string]struct{}) {
s.Mux.Lock()
defer s.Mux.Unlock()
s.routes = routes
}
// DeleteRoute removes a route from the network amp
func (s *State) DeleteRoute(network string) {
s.Mux.Lock()
defer s.Mux.Unlock()
delete(s.routes, network)
}
// GetRoutes return routes map
func (s *State) GetRoutes() map[string]struct{} {
s.Mux.RLock()
defer s.Mux.RUnlock()
return maps.Clone(s.routes)
}
@@ -1,4 +1,4 @@
package peer package status
import ( import (
"context" "context"
@@ -15,7 +15,6 @@ import (
"golang.org/x/exp/maps" "golang.org/x/exp/maps"
"google.golang.org/grpc/codes" "google.golang.org/grpc/codes"
gstatus "google.golang.org/grpc/status" gstatus "google.golang.org/grpc/status"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/timestamppb" "google.golang.org/protobuf/types/known/timestamppb"
firewall "github.com/netbirdio/netbird/client/firewall/manager" firewall "github.com/netbirdio/netbird/client/firewall/manager"
@@ -52,61 +51,6 @@ type RouterState struct {
Latency time.Duration Latency time.Duration
} }
// State contains the latest state of a peer
type State struct {
Mux *sync.RWMutex
IP string
IPv6 string
PubKey string
FQDN string
ConnStatus ConnStatus
ConnStatusUpdate time.Time
Relayed bool
LocalIceCandidateType string
RemoteIceCandidateType string
LocalIceCandidateEndpoint string
RemoteIceCandidateEndpoint string
RelayServerAddress string
LastWireguardHandshake time.Time
BytesTx int64
BytesRx int64
Latency time.Duration
RosenpassEnabled bool
SSHHostKey []byte
routes map[string]struct{}
}
// AddRoute add a single route to routes map
func (s *State) AddRoute(network string) {
s.Mux.Lock()
defer s.Mux.Unlock()
if s.routes == nil {
s.routes = make(map[string]struct{})
}
s.routes[network] = struct{}{}
}
// SetRoutes set state routes
func (s *State) SetRoutes(routes map[string]struct{}) {
s.Mux.Lock()
defer s.Mux.Unlock()
s.routes = routes
}
// DeleteRoute removes a route from the network amp
func (s *State) DeleteRoute(network string) {
s.Mux.Lock()
defer s.Mux.Unlock()
delete(s.routes, network)
}
// GetRoutes return routes map
func (s *State) GetRoutes() map[string]struct{} {
s.Mux.RLock()
defer s.Mux.RUnlock()
return maps.Clone(s.routes)
}
// LocalPeerState contains the latest state of the local peer // LocalPeerState contains the latest state of the local peer
type LocalPeerState struct { type LocalPeerState struct {
IP string IP string
@@ -154,20 +98,6 @@ type NSGroupState struct {
Error error Error error
} }
// FullStatus contains the full state held by the Status instance
type FullStatus struct {
Peers []State
ManagementState ManagementState
SignalState SignalState
LocalPeerState LocalPeerState
RosenpassState RosenpassState
Relays []relay.ProbeResult
NSGroupStates []NSGroupState
NumOfForwardingRules int
LazyConnectionEnabled bool
Events []*proto.SystemEvent
}
type StatusChangeSubscription struct { type StatusChangeSubscription struct {
peerID string peerID string
id string id string
@@ -189,11 +119,11 @@ func (s *StatusChangeSubscription) Events() chan map[string]RouterState {
return s.eventsChan return s.eventsChan
} }
// Status holds a state of peers, signal, management connections and relays. // Recorder holds a state of peers, signal, management connections and relays.
// mux is an RWMutex so hot read paths (notably PeerStateByIP, called for // mux is an RWMutex so hot read paths (notably PeerStateByIP, called for
// every private-service request) don't contend against each other. // every private-service request) don't contend against each other.
// Pure read methods take RLock; anything that mutates state takes Lock. // Pure read methods take RLock; anything that mutates state takes Lock.
type Status struct { type Recorder struct {
mux sync.RWMutex mux sync.RWMutex
muxRelays sync.RWMutex muxRelays sync.RWMutex
peers map[string]State peers map[string]State
@@ -254,9 +184,9 @@ type Status struct {
wgIface WGIfaceStatus wgIface WGIfaceStatus
} }
// NewRecorder returns a new Status instance // NewRecorder returns a new Recorder instance
func NewRecorder(mgmAddress string) *Status { func NewRecorder(mgmAddress string) *Recorder {
return &Status{ return &Recorder{
peers: make(map[string]State), peers: make(map[string]State),
ipToKey: make(map[string]string), ipToKey: make(map[string]string),
changeNotify: make(map[string]map[string]*StatusChangeSubscription), changeNotify: make(map[string]map[string]*StatusChangeSubscription),
@@ -270,20 +200,20 @@ func NewRecorder(mgmAddress string) *Status {
} }
} }
func (d *Status) SetRelayMgr(manager *relayClient.Manager) { func (d *Recorder) SetRelayMgr(manager *relayClient.Manager) {
d.muxRelays.Lock() d.muxRelays.Lock()
defer d.muxRelays.Unlock() defer d.muxRelays.Unlock()
d.relayMgr = manager d.relayMgr = manager
} }
func (d *Status) SetIngressGwMgr(ingressGwMgr *ingressgw.Manager) { func (d *Recorder) SetIngressGwMgr(ingressGwMgr *ingressgw.Manager) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
d.ingressGwMgr = ingressGwMgr d.ingressGwMgr = ingressGwMgr
} }
// ReplaceOfflinePeers replaces // ReplaceOfflinePeers replaces
func (d *Status) ReplaceOfflinePeers(replacement []State) { func (d *Recorder) ReplaceOfflinePeers(replacement []State) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
d.offlinePeers = make([]State, len(replacement)) d.offlinePeers = make([]State, len(replacement))
@@ -294,7 +224,7 @@ func (d *Status) ReplaceOfflinePeers(replacement []State) {
} }
// AddPeer adds peer to Daemon status map // AddPeer adds peer to Daemon status map
func (d *Status) AddPeer(peerPubKey string, fqdn string, ip string, ipv6 string) error { func (d *Recorder) AddPeer(peerPubKey string, fqdn string, ip string, ipv6 string) error {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -321,7 +251,7 @@ func (d *Status) AddPeer(peerPubKey string, fqdn string, ip string, ipv6 string)
} }
// GetPeer adds peer to Daemon status map // GetPeer adds peer to Daemon status map
func (d *Status) GetPeer(peerPubKey string) (State, error) { func (d *Recorder) GetPeer(peerPubKey string) (State, error) {
d.mux.RLock() d.mux.RLock()
defer d.mux.RUnlock() defer d.mux.RUnlock()
@@ -332,7 +262,7 @@ func (d *Status) GetPeer(peerPubKey string) (State, error) {
return state, nil return state, nil
} }
func (d *Status) PeerByIP(ip string) (string, bool) { func (d *Recorder) PeerByIP(ip string) (string, bool) {
d.mux.RLock() d.mux.RLock()
defer d.mux.RUnlock() defer d.mux.RUnlock()
@@ -349,7 +279,7 @@ func (d *Status) PeerByIP(ip string) (string, bool) {
// address so dual-stack peers are reachable on either family. Only // address so dual-stack peers are reachable on either family. Only
// active peers are matched; peers moved into the offline slice by // active peers are matched; peers moved into the offline slice by
// ReplaceOfflinePeers are intentionally treated as unknown. // ReplaceOfflinePeers are intentionally treated as unknown.
func (d *Status) PeerStateByIP(ip string) (State, bool) { func (d *Recorder) PeerStateByIP(ip string) (State, bool) {
if ip == "" { if ip == "" {
return State{}, false return State{}, false
} }
@@ -367,7 +297,7 @@ func (d *Status) PeerStateByIP(ip string) (State, bool) {
} }
// RemovePeer removes peer from Daemon status map // RemovePeer removes peer from Daemon status map
func (d *Status) RemovePeer(peerPubKey string) error { func (d *Recorder) RemovePeer(peerPubKey string) error {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -388,7 +318,7 @@ func (d *Status) RemovePeer(peerPubKey string) error {
} }
// UpdatePeerState updates peer status // UpdatePeerState updates peer status
func (d *Status) UpdatePeerState(receivedState State) error { func (d *Recorder) UpdatePeerState(receivedState State) error {
return d.updatePeer(receivedState.PubKey, return d.updatePeer(receivedState.PubKey,
func(_, updated State) bool { return updated.ConnStatus == StatusIdle }, func(_, updated State) bool { return updated.ConnStatus == StatusIdle },
func(peerState *State) { func(peerState *State) {
@@ -407,7 +337,7 @@ func (d *Status) UpdatePeerState(receivedState State) error {
}) })
} }
func (d *Status) AddPeerStateRoute(peer string, route string, resourceId route.ResID) error { func (d *Recorder) AddPeerStateRoute(peer string, route string, resourceId route.ResID) error {
d.mux.Lock() d.mux.Lock()
peerState, ok := d.peers[peer] peerState, ok := d.peers[peer]
@@ -433,7 +363,7 @@ func (d *Status) AddPeerStateRoute(peer string, route string, resourceId route.R
return nil return nil
} }
func (d *Status) RemovePeerStateRoute(peer string, route string) error { func (d *Recorder) RemovePeerStateRoute(peer string, route string) error {
d.mux.Lock() d.mux.Lock()
peerState, ok := d.peers[peer] peerState, ok := d.peers[peer]
@@ -461,7 +391,7 @@ func (d *Status) RemovePeerStateRoute(peer string, route string) error {
// CheckRoutes checks if the source and destination addresses are within the same route // CheckRoutes checks if the source and destination addresses are within the same route
// and returns the resource ID of the route that contains the addresses // and returns the resource ID of the route that contains the addresses
func (d *Status) CheckRoutes(ip netip.Addr) ([]byte, bool) { func (d *Recorder) CheckRoutes(ip netip.Addr) ([]byte, bool) {
if d == nil { if d == nil {
return nil, false return nil, false
} }
@@ -469,7 +399,7 @@ func (d *Status) CheckRoutes(ip netip.Addr) ([]byte, bool) {
return []byte(resId), isExitNode return []byte(resId), isExitNode
} }
func (d *Status) UpdatePeerICEState(receivedState State) error { func (d *Recorder) UpdatePeerICEState(receivedState State) error {
return d.updatePeer(receivedState.PubKey, hasStatusOrRelayedChange, func(peerState *State) { return d.updatePeer(receivedState.PubKey, hasStatusOrRelayedChange, func(peerState *State) {
peerState.ConnStatus = receivedState.ConnStatus peerState.ConnStatus = receivedState.ConnStatus
peerState.ConnStatusUpdate = receivedState.ConnStatusUpdate peerState.ConnStatusUpdate = receivedState.ConnStatusUpdate
@@ -482,7 +412,7 @@ func (d *Status) UpdatePeerICEState(receivedState State) error {
}) })
} }
func (d *Status) UpdatePeerRelayedState(receivedState State) error { func (d *Recorder) UpdatePeerRelayedState(receivedState State) error {
return d.updatePeer(receivedState.PubKey, hasStatusOrRelayedChange, func(peerState *State) { return d.updatePeer(receivedState.PubKey, hasStatusOrRelayedChange, func(peerState *State) {
peerState.ConnStatus = receivedState.ConnStatus peerState.ConnStatus = receivedState.ConnStatus
peerState.ConnStatusUpdate = receivedState.ConnStatusUpdate peerState.ConnStatusUpdate = receivedState.ConnStatusUpdate
@@ -492,7 +422,7 @@ func (d *Status) UpdatePeerRelayedState(receivedState State) error {
}) })
} }
func (d *Status) UpdatePeerRelayedStateToDisconnected(receivedState State) error { func (d *Recorder) UpdatePeerRelayedStateToDisconnected(receivedState State) error {
return d.updatePeer(receivedState.PubKey, hasStatusOrRelayedChange, func(peerState *State) { return d.updatePeer(receivedState.PubKey, hasStatusOrRelayedChange, func(peerState *State) {
peerState.ConnStatus = receivedState.ConnStatus peerState.ConnStatus = receivedState.ConnStatus
peerState.Relayed = receivedState.Relayed peerState.Relayed = receivedState.Relayed
@@ -501,7 +431,7 @@ func (d *Status) UpdatePeerRelayedStateToDisconnected(receivedState State) error
}) })
} }
func (d *Status) UpdatePeerICEStateToDisconnected(receivedState State) error { func (d *Recorder) UpdatePeerICEStateToDisconnected(receivedState State) error {
return d.updatePeer(receivedState.PubKey, hasStatusOrRelayedChange, func(peerState *State) { return d.updatePeer(receivedState.PubKey, hasStatusOrRelayedChange, func(peerState *State) {
peerState.ConnStatus = receivedState.ConnStatus peerState.ConnStatus = receivedState.ConnStatus
peerState.Relayed = receivedState.Relayed peerState.Relayed = receivedState.Relayed
@@ -514,7 +444,7 @@ func (d *Status) UpdatePeerICEStateToDisconnected(receivedState State) error {
} }
// UpdateWireGuardPeerState updates the WireGuard bits of the peer state // UpdateWireGuardPeerState updates the WireGuard bits of the peer state
func (d *Status) UpdateWireGuardPeerState(pubKey string, wgStats configurer.WGStats) error { func (d *Recorder) UpdateWireGuardPeerState(pubKey string, wgStats configurer.WGStats) error {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -534,7 +464,7 @@ func (d *Status) UpdateWireGuardPeerState(pubKey string, wgStats configurer.WGSt
// updatePeer applies mutate to the stored peer state and runs the list, router // updatePeer applies mutate to the stored peer state and runs the list, router
// and state-change notifications outside the lock // and state-change notifications outside the lock
func (d *Status) updatePeer(pubKey string, notifyRouter func(old, updated State) bool, mutate func(*State)) error { func (d *Recorder) updatePeer(pubKey string, notifyRouter func(old, updated State) bool, mutate func(*State)) error {
d.mux.Lock() d.mux.Lock()
peerState, ok := d.peers[pubKey] peerState, ok := d.peers[pubKey]
@@ -564,12 +494,8 @@ func (d *Status) updatePeer(pubKey string, notifyRouter func(old, updated State)
return nil return nil
} }
func hasStatusOrRelayedChange(old, updated State) bool {
return old.Relayed != updated.Relayed || old.ConnStatus != updated.ConnStatus
}
// UpdatePeerFQDN update peer's state fqdn only // UpdatePeerFQDN update peer's state fqdn only
func (d *Status) UpdatePeerFQDN(peerPubKey, fqdn string) error { func (d *Recorder) UpdatePeerFQDN(peerPubKey, fqdn string) error {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -585,7 +511,7 @@ func (d *Status) UpdatePeerFQDN(peerPubKey, fqdn string) error {
} }
// UpdatePeerSSHHostKey updates peer's SSH host key // UpdatePeerSSHHostKey updates peer's SSH host key
func (d *Status) UpdatePeerSSHHostKey(peerPubKey string, sshHostKey []byte) error { func (d *Recorder) UpdatePeerSSHHostKey(peerPubKey string, sshHostKey []byte) error {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -601,7 +527,7 @@ func (d *Status) UpdatePeerSSHHostKey(peerPubKey string, sshHostKey []byte) erro
} }
// FinishPeerListModifications this event invoke the notification // FinishPeerListModifications this event invoke the notification
func (d *Status) FinishPeerListModifications() { func (d *Recorder) FinishPeerListModifications() {
d.mux.Lock() d.mux.Lock()
if !d.peerListChangedForNotification { if !d.peerListChangedForNotification {
@@ -634,7 +560,7 @@ func (d *Status) FinishPeerListModifications() {
d.notifyStateChange() d.notifyStateChange()
} }
func (d *Status) SubscribeToPeerStateChanges(ctx context.Context, peerID string) *StatusChangeSubscription { func (d *Recorder) SubscribeToPeerStateChanges(ctx context.Context, peerID string) *StatusChangeSubscription {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -647,7 +573,7 @@ func (d *Status) SubscribeToPeerStateChanges(ctx context.Context, peerID string)
return sub return sub
} }
func (d *Status) UnsubscribePeerStateChanges(subscription *StatusChangeSubscription) { func (d *Recorder) UnsubscribePeerStateChanges(subscription *StatusChangeSubscription) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -672,14 +598,14 @@ func (d *Status) UnsubscribePeerStateChanges(subscription *StatusChangeSubscript
} }
// GetLocalPeerState returns the local peer state // GetLocalPeerState returns the local peer state
func (d *Status) GetLocalPeerState() LocalPeerState { func (d *Recorder) GetLocalPeerState() LocalPeerState {
d.mux.RLock() d.mux.RLock()
defer d.mux.RUnlock() defer d.mux.RUnlock()
return d.localPeer.Clone() return d.localPeer.Clone()
} }
// UpdateLocalPeerState updates local peer status // UpdateLocalPeerState updates local peer status
func (d *Status) UpdateLocalPeerState(localPeerState LocalPeerState) { func (d *Recorder) UpdateLocalPeerState(localPeerState LocalPeerState) {
d.mux.Lock() d.mux.Lock()
d.localPeer = localPeerState d.localPeer = localPeerState
fqdn := d.localPeer.FQDN fqdn := d.localPeer.FQDN
@@ -699,7 +625,7 @@ func (d *Status) UpdateLocalPeerState(localPeerState LocalPeerState) {
// disabled or the peer is not SSO-tracked). Same-value updates are no-ops; // disabled or the peer is not SSO-tracked). Same-value updates are no-ops;
// real changes fan out via notifyStateChange so SubscribeStatus consumers // real changes fan out via notifyStateChange so SubscribeStatus consumers
// pick up the new deadline on their next read. // pick up the new deadline on their next read.
func (d *Status) SetSessionExpiresAt(deadline time.Time) { func (d *Recorder) SetSessionExpiresAt(deadline time.Time) {
d.mux.Lock() d.mux.Lock()
if d.sessionExpiresAt.Equal(deadline) { if d.sessionExpiresAt.Equal(deadline) {
d.mux.Unlock() d.mux.Unlock()
@@ -716,14 +642,14 @@ func (d *Status) SetSessionExpiresAt(deadline time.Time) {
// CLI status) render it as "expired" rather than hiding it — masking it as // CLI status) render it as "expired" rather than hiding it — masking it as
// "none" would blank the UI at the exact moment it should say the session // "none" would blank the UI at the exact moment it should say the session
// ended. // ended.
func (d *Status) GetSessionExpiresAt() time.Time { func (d *Recorder) GetSessionExpiresAt() time.Time {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
return d.sessionExpiresAt return d.sessionExpiresAt
} }
// AddLocalPeerStateRoute adds a route to the local peer state // AddLocalPeerStateRoute adds a route to the local peer state
func (d *Status) AddLocalPeerStateRoute(route string, resourceId route.ResID) { func (d *Recorder) AddLocalPeerStateRoute(route string, resourceId route.ResID) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -740,7 +666,7 @@ func (d *Status) AddLocalPeerStateRoute(route string, resourceId route.ResID) {
} }
// RemoveLocalPeerStateRoute removes a route from the local peer state // RemoveLocalPeerStateRoute removes a route from the local peer state
func (d *Status) RemoveLocalPeerStateRoute(route string) { func (d *Recorder) RemoveLocalPeerStateRoute(route string) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -753,7 +679,7 @@ func (d *Status) RemoveLocalPeerStateRoute(route string) {
} }
// AddResolvedIPLookupEntry adds a resolved IP lookup entry // AddResolvedIPLookupEntry adds a resolved IP lookup entry
func (d *Status) AddResolvedIPLookupEntry(prefix netip.Prefix, resourceId route.ResID) { func (d *Recorder) AddResolvedIPLookupEntry(prefix netip.Prefix, resourceId route.ResID) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -761,7 +687,7 @@ func (d *Status) AddResolvedIPLookupEntry(prefix netip.Prefix, resourceId route.
} }
// RemoveResolvedIPLookupEntry removes a resolved IP lookup entry // RemoveResolvedIPLookupEntry removes a resolved IP lookup entry
func (d *Status) RemoveResolvedIPLookupEntry(route string) { func (d *Recorder) RemoveResolvedIPLookupEntry(route string) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -772,7 +698,7 @@ func (d *Status) RemoveResolvedIPLookupEntry(route string) {
} }
// CleanLocalPeerStateRoutes cleans all routes from the local peer state // CleanLocalPeerStateRoutes cleans all routes from the local peer state
func (d *Status) CleanLocalPeerStateRoutes() { func (d *Recorder) CleanLocalPeerStateRoutes() {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -780,7 +706,7 @@ func (d *Status) CleanLocalPeerStateRoutes() {
} }
// CleanLocalPeerState cleans local peer status // CleanLocalPeerState cleans local peer status
func (d *Status) CleanLocalPeerState() { func (d *Recorder) CleanLocalPeerState() {
d.mux.Lock() d.mux.Lock()
d.localPeer = LocalPeerState{} d.localPeer = LocalPeerState{}
fqdn := d.localPeer.FQDN fqdn := d.localPeer.FQDN
@@ -792,7 +718,7 @@ func (d *Status) CleanLocalPeerState() {
} }
// MarkManagementDisconnected sets ManagementState to disconnected // MarkManagementDisconnected sets ManagementState to disconnected
func (d *Status) MarkManagementDisconnected(err error) { func (d *Recorder) MarkManagementDisconnected(err error) {
d.mux.Lock() d.mux.Lock()
// Health checks re-mark the same state on every probe; skip the fan-out // Health checks re-mark the same state on every probe; skip the fan-out
// when nothing actually changed so we don't flood SubscribeStatus // when nothing actually changed so we don't flood SubscribeStatus
@@ -812,7 +738,7 @@ func (d *Status) MarkManagementDisconnected(err error) {
} }
// MarkManagementConnected sets ManagementState to connected // MarkManagementConnected sets ManagementState to connected
func (d *Status) MarkManagementConnected() { func (d *Recorder) MarkManagementConnected() {
d.mux.Lock() d.mux.Lock()
if d.managementState && d.managementError == nil { if d.managementState && d.managementError == nil {
d.mux.Unlock() d.mux.Unlock()
@@ -829,35 +755,35 @@ func (d *Status) MarkManagementConnected() {
} }
// UpdateSignalAddress update the address of the signal server // UpdateSignalAddress update the address of the signal server
func (d *Status) UpdateSignalAddress(signalURL string) { func (d *Recorder) UpdateSignalAddress(signalURL string) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
d.signalAddress = signalURL d.signalAddress = signalURL
} }
// UpdateManagementAddress update the address of the management server // UpdateManagementAddress update the address of the management server
func (d *Status) UpdateManagementAddress(mgmAddress string) { func (d *Recorder) UpdateManagementAddress(mgmAddress string) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
d.mgmAddress = mgmAddress d.mgmAddress = mgmAddress
} }
// UpdateRosenpass update the Rosenpass configuration // UpdateRosenpass update the Rosenpass configuration
func (d *Status) UpdateRosenpass(rosenpassEnabled, rosenpassPermissive bool) { func (d *Recorder) UpdateRosenpass(rosenpassEnabled, rosenpassPermissive bool) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
d.rosenpassPermissive = rosenpassPermissive d.rosenpassPermissive = rosenpassPermissive
d.rosenpassEnabled = rosenpassEnabled d.rosenpassEnabled = rosenpassEnabled
} }
func (d *Status) UpdateLazyConnection(enabled bool) { func (d *Recorder) UpdateLazyConnection(enabled bool) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
d.lazyConnectionEnabled = enabled d.lazyConnectionEnabled = enabled
} }
// MarkSignalDisconnected sets SignalState to disconnected // MarkSignalDisconnected sets SignalState to disconnected
func (d *Status) MarkSignalDisconnected(err error) { func (d *Recorder) MarkSignalDisconnected(err error) {
d.mux.Lock() d.mux.Lock()
if !d.signalState && errors.Is(d.signalError, err) { if !d.signalState && errors.Is(d.signalError, err) {
d.mux.Unlock() d.mux.Unlock()
@@ -874,7 +800,7 @@ func (d *Status) MarkSignalDisconnected(err error) {
} }
// MarkSignalConnected sets SignalState to connected // MarkSignalConnected sets SignalState to connected
func (d *Status) MarkSignalConnected() { func (d *Recorder) MarkSignalConnected() {
d.mux.Lock() d.mux.Lock()
if d.signalState && d.signalError == nil { if d.signalState && d.signalError == nil {
d.mux.Unlock() d.mux.Unlock()
@@ -890,19 +816,19 @@ func (d *Status) MarkSignalConnected() {
d.notifyStateChange() d.notifyStateChange()
} }
func (d *Status) UpdateRelayStates(relayResults []relay.ProbeResult) { func (d *Recorder) UpdateRelayStates(relayResults []relay.ProbeResult) {
d.muxRelays.Lock() d.muxRelays.Lock()
defer d.muxRelays.Unlock() defer d.muxRelays.Unlock()
d.relayStates = relayResults d.relayStates = relayResults
} }
func (d *Status) UpdateDNSStates(dnsStates []NSGroupState) { func (d *Recorder) UpdateDNSStates(dnsStates []NSGroupState) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
d.nsGroupStates = dnsStates d.nsGroupStates = dnsStates
} }
func (d *Status) UpdateResolvedDomainsStates(originalDomain domain.Domain, resolvedDomain domain.Domain, prefixes []netip.Prefix, resourceId route.ResID) { func (d *Recorder) UpdateResolvedDomainsStates(originalDomain domain.Domain, resolvedDomain domain.Domain, prefixes []netip.Prefix, resourceId route.ResID) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -917,7 +843,7 @@ func (d *Status) UpdateResolvedDomainsStates(originalDomain domain.Domain, resol
} }
} }
func (d *Status) DeleteResolvedDomainsStates(domain domain.Domain) { func (d *Recorder) DeleteResolvedDomainsStates(domain domain.Domain) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -933,7 +859,7 @@ func (d *Status) DeleteResolvedDomainsStates(domain domain.Domain) {
} }
} }
func (d *Status) GetRosenpassState() RosenpassState { func (d *Recorder) GetRosenpassState() RosenpassState {
d.mux.RLock() d.mux.RLock()
defer d.mux.RUnlock() defer d.mux.RUnlock()
return RosenpassState{ return RosenpassState{
@@ -942,13 +868,13 @@ func (d *Status) GetRosenpassState() RosenpassState {
} }
} }
func (d *Status) GetLazyConnection() bool { func (d *Recorder) GetLazyConnection() bool {
d.mux.RLock() d.mux.RLock()
defer d.mux.RUnlock() defer d.mux.RUnlock()
return d.lazyConnectionEnabled return d.lazyConnectionEnabled
} }
func (d *Status) GetManagementState() ManagementState { func (d *Recorder) GetManagementState() ManagementState {
d.mux.RLock() d.mux.RLock()
defer d.mux.RUnlock() defer d.mux.RUnlock()
return ManagementState{ return ManagementState{
@@ -958,7 +884,7 @@ func (d *Status) GetManagementState() ManagementState {
} }
} }
func (d *Status) UpdateLatency(pubKey string, latency time.Duration) error { func (d *Recorder) UpdateLatency(pubKey string, latency time.Duration) error {
if latency <= 0 { if latency <= 0 {
return nil return nil
} }
@@ -975,7 +901,7 @@ func (d *Status) UpdateLatency(pubKey string, latency time.Duration) error {
} }
// IsLoginRequired determines if a peer's login has expired. // IsLoginRequired determines if a peer's login has expired.
func (d *Status) IsLoginRequired() bool { func (d *Recorder) IsLoginRequired() bool {
d.mux.RLock() d.mux.RLock()
defer d.mux.RUnlock() defer d.mux.RUnlock()
@@ -991,7 +917,7 @@ func (d *Status) IsLoginRequired() bool {
return false return false
} }
func (d *Status) GetSignalState() SignalState { func (d *Recorder) GetSignalState() SignalState {
d.mux.RLock() d.mux.RLock()
defer d.mux.RUnlock() defer d.mux.RUnlock()
return SignalState{ return SignalState{
@@ -1002,7 +928,7 @@ func (d *Status) GetSignalState() SignalState {
} }
// GetRelayStates returns the stun/turn/permanent relay states // GetRelayStates returns the stun/turn/permanent relay states
func (d *Status) GetRelayStates() []relay.ProbeResult { func (d *Recorder) GetRelayStates() []relay.ProbeResult {
d.muxRelays.RLock() d.muxRelays.RLock()
if d.relayMgr == nil { if d.relayMgr == nil {
defer d.muxRelays.RUnlock() defer d.muxRelays.RUnlock()
@@ -1041,7 +967,7 @@ func (d *Status) GetRelayStates() []relay.ProbeResult {
return relayStates return relayStates
} }
func (d *Status) ForwardingRules() []firewall.ForwardRule { func (d *Recorder) ForwardingRules() []firewall.ForwardRule {
d.mux.RLock() d.mux.RLock()
defer d.mux.RUnlock() defer d.mux.RUnlock()
if d.ingressGwMgr == nil { if d.ingressGwMgr == nil {
@@ -1051,7 +977,7 @@ func (d *Status) ForwardingRules() []firewall.ForwardRule {
return d.ingressGwMgr.Rules() return d.ingressGwMgr.Rules()
} }
func (d *Status) GetDNSStates() []NSGroupState { func (d *Recorder) GetDNSStates() []NSGroupState {
d.mux.RLock() d.mux.RLock()
defer d.mux.RUnlock() defer d.mux.RUnlock()
@@ -1059,14 +985,14 @@ func (d *Status) GetDNSStates() []NSGroupState {
return slices.Clone(d.nsGroupStates) return slices.Clone(d.nsGroupStates)
} }
func (d *Status) GetResolvedDomainsStates() map[domain.Domain]ResolvedDomainInfo { func (d *Recorder) GetResolvedDomainsStates() map[domain.Domain]ResolvedDomainInfo {
d.mux.RLock() d.mux.RLock()
defer d.mux.RUnlock() defer d.mux.RUnlock()
return maps.Clone(d.resolvedDomainsStates) return maps.Clone(d.resolvedDomainsStates)
} }
// GetFullStatus gets full status // GetFullStatus gets full status
func (d *Status) GetFullStatus() FullStatus { func (d *Recorder) GetFullStatus() FullStatus {
fullStatus := FullStatus{ fullStatus := FullStatus{
ManagementState: d.GetManagementState(), ManagementState: d.GetManagementState(),
SignalState: d.GetSignalState(), SignalState: d.GetSignalState(),
@@ -1092,30 +1018,30 @@ func (d *Status) GetFullStatus() FullStatus {
} }
// ClientStart will notify all listeners about the new service state // ClientStart will notify all listeners about the new service state
func (d *Status) ClientStart() { func (d *Recorder) ClientStart() {
d.notifier.clientStart() d.notifier.clientStart()
d.notifyStateChange() d.notifyStateChange()
} }
// ClientStop will notify all listeners about the new service state // ClientStop will notify all listeners about the new service state
func (d *Status) ClientStop() { func (d *Recorder) ClientStop() {
d.notifier.clientStop() d.notifier.clientStop()
d.notifyStateChange() d.notifyStateChange()
} }
// ClientTeardown will notify all listeners about the service is under teardown // ClientTeardown will notify all listeners about the service is under teardown
func (d *Status) ClientTeardown() { func (d *Recorder) ClientTeardown() {
d.notifier.clientTearDown() d.notifier.clientTearDown()
d.notifyStateChange() d.notifyStateChange()
} }
// SetConnectionListener set a listener to the notifier // SetConnectionListener set a listener to the notifier
func (d *Status) SetConnectionListener(listener Listener) { func (d *Recorder) SetConnectionListener(listener Listener) {
d.notifier.setListener(listener) d.notifier.setListener(listener)
} }
// RemoveConnectionListener remove the listener from the notifier // RemoveConnectionListener remove the listener from the notifier
func (d *Status) RemoveConnectionListener() { func (d *Recorder) RemoveConnectionListener() {
d.notifier.removeListener() d.notifier.removeListener()
} }
@@ -1123,7 +1049,7 @@ func (d *Status) RemoveConnectionListener() {
// Caller MUST hold d.mux. Returns nil when there are no subscribers for peerID // Caller MUST hold d.mux. Returns nil when there are no subscribers for peerID
// or when notify is false. The snapshot is consumed later by dispatchRouterPeers // or when notify is false. The snapshot is consumed later by dispatchRouterPeers
// outside the lock so the channel send cannot stall any d.mux holder. // outside the lock so the channel send cannot stall any d.mux holder.
func (d *Status) snapshotRouterPeersLocked(peerID string, notify bool) map[string]RouterState { func (d *Recorder) snapshotRouterPeersLocked(peerID string, notify bool) map[string]RouterState {
if !notify { if !notify {
return nil return nil
} }
@@ -1152,7 +1078,7 @@ func (d *Status) snapshotRouterPeersLocked(peerID string, notify bool) map[strin
// channels, then sends outside the lock so a slow consumer cannot block other // channels, then sends outside the lock so a slow consumer cannot block other
// d.mux holders. The send itself stays blocking (only short-circuited by the // d.mux holders. The send itself stays blocking (only short-circuited by the
// subscriber's context) so peer state transitions are not silently dropped. // subscriber's context) so peer state transitions are not silently dropped.
func (d *Status) dispatchRouterPeers(peerID string, routerPeers map[string]RouterState) { func (d *Recorder) dispatchRouterPeers(peerID string, routerPeers map[string]RouterState) {
if routerPeers == nil { if routerPeers == nil {
return return
} }
@@ -1175,12 +1101,12 @@ func (d *Status) dispatchRouterPeers(peerID string, routerPeers map[string]Route
} }
} }
func (d *Status) numOfPeers() int { func (d *Recorder) numOfPeers() int {
return len(d.peers) + len(d.offlinePeers) return len(d.peers) + len(d.offlinePeers)
} }
// PublishEvent adds an event to the queue and distributes it to all subscribers // PublishEvent adds an event to the queue and distributes it to all subscribers
func (d *Status) PublishEvent( func (d *Recorder) PublishEvent(
severity proto.SystemEvent_Severity, severity proto.SystemEvent_Severity,
category proto.SystemEvent_Category, category proto.SystemEvent_Category,
msg string, msg string,
@@ -1214,7 +1140,7 @@ func (d *Status) PublishEvent(
} }
// SubscribeToEvents returns a new event subscription // SubscribeToEvents returns a new event subscription
func (d *Status) SubscribeToEvents() *EventSubscription { func (d *Recorder) SubscribeToEvents() *EventSubscription {
d.eventMux.Lock() d.eventMux.Lock()
defer d.eventMux.Unlock() defer d.eventMux.Unlock()
@@ -1229,7 +1155,7 @@ func (d *Status) SubscribeToEvents() *EventSubscription {
} }
// UnsubscribeFromEvents removes an event subscription // UnsubscribeFromEvents removes an event subscription
func (d *Status) UnsubscribeFromEvents(sub *EventSubscription) { func (d *Recorder) UnsubscribeFromEvents(sub *EventSubscription) {
if sub == nil { if sub == nil {
return return
} }
@@ -1244,7 +1170,7 @@ func (d *Status) UnsubscribeFromEvents(sub *EventSubscription) {
} }
// GetEventHistory returns all events in the queue // GetEventHistory returns all events in the queue
func (d *Status) GetEventHistory() []*proto.SystemEvent { func (d *Recorder) GetEventHistory() []*proto.SystemEvent {
return d.eventQueue.GetAll() return d.eventQueue.GetAll()
} }
@@ -1253,7 +1179,7 @@ func (d *Status) GetEventHistory() []*proto.SystemEvent {
// address change / peers-list change). The channel is buffered to one // address change / peers-list change). The channel is buffered to one
// pending tick so a coalesced burst still wakes the consumer exactly // pending tick so a coalesced burst still wakes the consumer exactly
// once. Pass the returned id to UnsubscribeFromStateChanges to detach. // once. Pass the returned id to UnsubscribeFromStateChanges to detach.
func (d *Status) SubscribeToStateChanges() (string, <-chan struct{}) { func (d *Recorder) SubscribeToStateChanges() (string, <-chan struct{}) {
d.stateChangeMux.Lock() d.stateChangeMux.Lock()
defer d.stateChangeMux.Unlock() defer d.stateChangeMux.Unlock()
@@ -1266,7 +1192,7 @@ func (d *Status) SubscribeToStateChanges() (string, <-chan struct{}) {
// UnsubscribeFromStateChanges releases a SubscribeToStateChanges channel // UnsubscribeFromStateChanges releases a SubscribeToStateChanges channel
// and closes it so any consumer goroutine selecting on the channel // and closes it so any consumer goroutine selecting on the channel
// unblocks cleanly. // unblocks cleanly.
func (d *Status) UnsubscribeFromStateChanges(id string) { func (d *Recorder) UnsubscribeFromStateChanges(id string) {
d.stateChangeMux.Lock() d.stateChangeMux.Lock()
defer d.stateChangeMux.Unlock() defer d.stateChangeMux.Unlock()
@@ -1280,7 +1206,7 @@ func (d *Status) UnsubscribeFromStateChanges(id string) {
// the tick if a subscriber's buffer is full — by definition the consumer // the tick if a subscriber's buffer is full — by definition the consumer
// is already going to fetch the latest snapshot, so multiple pending ticks // is already going to fetch the latest snapshot, so multiple pending ticks
// would be redundant. // would be redundant.
func (d *Status) notifyStateChange() { func (d *Recorder) notifyStateChange() {
d.stateChangeMux.Lock() d.stateChangeMux.Lock()
defer d.stateChangeMux.Unlock() defer d.stateChangeMux.Unlock()
@@ -1300,7 +1226,7 @@ func (d *Status) notifyStateChange() {
// the previous snapshot until an unrelated peer/management/signal // the previous snapshot until an unrelated peer/management/signal
// change happens to fire notifyStateChange, leaving the UI's status // change happens to fire notifyStateChange, leaving the UI's status
// out of sync with the daemon. // out of sync with the daemon.
func (d *Status) NotifyStateChange() { func (d *Recorder) NotifyStateChange() {
d.notifyStateChange() d.notifyStateChange()
} }
@@ -1309,7 +1235,7 @@ func (d *Status) NotifyStateChange() {
// changes the available routes or when a selection is applied — the peer // changes the available routes or when a selection is applied — the peer
// status itself only records actively-routed (chosen) networks, so without // status itself only records actively-routed (chosen) networks, so without
// this bump a candidate route appearing/disappearing would never reach the UI. // this bump a candidate route appearing/disappearing would never reach the UI.
func (d *Status) BumpNetworksRevision() { func (d *Recorder) BumpNetworksRevision() {
d.networksRevision.Add(1) d.networksRevision.Add(1)
d.notifyStateChange() d.notifyStateChange()
} }
@@ -1317,18 +1243,18 @@ func (d *Status) BumpNetworksRevision() {
// GetNetworksRevision returns the current routed-networks revision, surfaced in // GetNetworksRevision returns the current routed-networks revision, surfaced in
// the status snapshot so the UI can detect route/selection changes (see // the status snapshot so the UI can detect route/selection changes (see
// BumpNetworksRevision). // BumpNetworksRevision).
func (d *Status) GetNetworksRevision() uint64 { func (d *Recorder) GetNetworksRevision() uint64 {
return d.networksRevision.Load() return d.networksRevision.Load()
} }
func (d *Status) SetWgIface(wgInterface WGIfaceStatus) { func (d *Recorder) SetWgIface(wgInterface WGIfaceStatus) {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
d.wgIface = wgInterface d.wgIface = wgInterface
} }
func (d *Status) PeersStatus() (*configurer.Stats, error) { func (d *Recorder) PeersStatus() (*configurer.Stats, error) {
d.mux.RLock() d.mux.RLock()
defer d.mux.RUnlock() defer d.mux.RUnlock()
if d.wgIface == nil { if d.wgIface == nil {
@@ -1341,7 +1267,7 @@ func (d *Status) PeersStatus() (*configurer.Stats, error) {
// RefreshWireGuardStats fetches fresh WireGuard statistics from the interface // RefreshWireGuardStats fetches fresh WireGuard statistics from the interface
// and updates the cached peer states. This ensures accurate handshake times and // and updates the cached peer states. This ensures accurate handshake times and
// transfer statistics in status reports without running full health probes. // transfer statistics in status reports without running full health probes.
func (d *Status) RefreshWireGuardStats() error { func (d *Recorder) RefreshWireGuardStats() error {
d.mux.Lock() d.mux.Lock()
defer d.mux.Unlock() defer d.mux.Unlock()
@@ -1370,140 +1296,6 @@ func (d *Status) RefreshWireGuardStats() error {
return nil return nil
} }
type EventQueue struct { func hasStatusOrRelayedChange(old, updated State) bool {
maxSize int return old.Relayed != updated.Relayed || old.ConnStatus != updated.ConnStatus
events []*proto.SystemEvent
mutex sync.RWMutex
}
func NewEventQueue(size int) *EventQueue {
return &EventQueue{
maxSize: size,
events: make([]*proto.SystemEvent, 0, size),
}
}
func (q *EventQueue) Add(event *proto.SystemEvent) {
q.mutex.Lock()
defer q.mutex.Unlock()
q.events = append(q.events, event)
if len(q.events) > q.maxSize {
q.events = q.events[len(q.events)-q.maxSize:]
}
}
func (q *EventQueue) GetAll() []*proto.SystemEvent {
q.mutex.RLock()
defer q.mutex.RUnlock()
return slices.Clone(q.events)
}
type EventSubscription struct {
id string
events chan *proto.SystemEvent
}
func (s *EventSubscription) Events() <-chan *proto.SystemEvent {
return s.events
}
// ToProto converts FullStatus to proto.FullStatus.
func (fs FullStatus) ToProto() *proto.FullStatus {
pbFullStatus := proto.FullStatus{
ManagementState: &proto.ManagementState{},
SignalState: &proto.SignalState{},
LocalPeerState: &proto.LocalPeerState{},
Peers: []*proto.PeerState{},
}
pbFullStatus.ManagementState.URL = fs.ManagementState.URL
pbFullStatus.ManagementState.Connected = fs.ManagementState.Connected
if err := fs.ManagementState.Error; err != nil {
pbFullStatus.ManagementState.Error = err.Error()
}
pbFullStatus.SignalState.URL = fs.SignalState.URL
pbFullStatus.SignalState.Connected = fs.SignalState.Connected
if err := fs.SignalState.Error; err != nil {
pbFullStatus.SignalState.Error = err.Error()
}
pbFullStatus.LocalPeerState.IP = fs.LocalPeerState.IP
pbFullStatus.LocalPeerState.Ipv6 = fs.LocalPeerState.IPv6
pbFullStatus.LocalPeerState.PubKey = fs.LocalPeerState.PubKey
pbFullStatus.LocalPeerState.KernelInterface = fs.LocalPeerState.KernelInterface
pbFullStatus.LocalPeerState.Fqdn = fs.LocalPeerState.FQDN
pbFullStatus.LocalPeerState.WgPort = int32(fs.LocalPeerState.WgPort)
pbFullStatus.LocalPeerState.RosenpassPermissive = fs.RosenpassState.Permissive
pbFullStatus.LocalPeerState.RosenpassEnabled = fs.RosenpassState.Enabled
pbFullStatus.NumberOfForwardingRules = int32(fs.NumOfForwardingRules)
pbFullStatus.LazyConnectionEnabled = fs.LazyConnectionEnabled
pbFullStatus.LocalPeerState.Networks = maps.Keys(fs.LocalPeerState.Routes)
for _, peerState := range fs.Peers {
networks := maps.Keys(peerState.GetRoutes())
pbPeerState := &proto.PeerState{
IP: peerState.IP,
Ipv6: peerState.IPv6,
PubKey: peerState.PubKey,
ConnStatus: peerState.ConnStatus.String(),
ConnStatusUpdate: timestamppb.New(peerState.ConnStatusUpdate),
Relayed: peerState.Relayed,
LocalIceCandidateType: peerState.LocalIceCandidateType,
RemoteIceCandidateType: peerState.RemoteIceCandidateType,
LocalIceCandidateEndpoint: peerState.LocalIceCandidateEndpoint,
RemoteIceCandidateEndpoint: peerState.RemoteIceCandidateEndpoint,
RelayAddress: peerState.RelayServerAddress,
Fqdn: peerState.FQDN,
LastWireguardHandshake: timestamppb.New(peerState.LastWireguardHandshake),
BytesRx: peerState.BytesRx,
BytesTx: peerState.BytesTx,
RosenpassEnabled: peerState.RosenpassEnabled,
Networks: networks,
Latency: durationpb.New(peerState.Latency),
SshHostKey: peerState.SSHHostKey,
}
pbFullStatus.Peers = append(pbFullStatus.Peers, pbPeerState)
}
for _, relayState := range fs.Relays {
pbRelayState := &proto.RelayState{
URI: relayState.URI,
Available: relayState.Err == nil,
Transport: relayState.Transport,
}
if err := relayState.Err; err != nil {
pbRelayState.Error = err.Error()
}
pbFullStatus.Relays = append(pbFullStatus.Relays, pbRelayState)
}
for _, dnsState := range fs.NSGroupStates {
var err string
if dnsState.Error != nil {
err = dnsState.Error.Error()
}
var servers []string
for _, server := range dnsState.Servers {
servers = append(servers, server.String())
}
pbDnsState := &proto.NSGroupState{
Servers: servers,
Domains: dnsState.Domains,
Enabled: dnsState.Enabled,
Error: err,
}
pbFullStatus.DnsServers = append(pbFullStatus.DnsServers, pbDnsState)
}
pbFullStatus.Events = fs.Events
return &pbFullStatus
} }
@@ -1,4 +1,4 @@
package peer package status
import ( import (
"context" "context"
@@ -1,4 +1,4 @@
package peer package status
import ( import (
"net/netip" "net/netip"
+36
View File
@@ -0,0 +1,36 @@
package peer
import "github.com/netbirdio/netbird/client/internal/peer/status"
// Transitional aliases re-exporting the peer status recorder from its own
// package. Callers are being migrated to reference the status package
// directly; these aliases will be removed once the migration completes.
type (
Status = status.Recorder
State = status.State
ConnStatus = status.ConnStatus
FullStatus = status.FullStatus
RouterState = status.RouterState
LocalPeerState = status.LocalPeerState
SignalState = status.SignalState
ManagementState = status.ManagementState
RosenpassState = status.RosenpassState
NSGroupState = status.NSGroupState
ResolvedDomainInfo = status.ResolvedDomainInfo
StatusChangeSubscription = status.StatusChangeSubscription
EventQueue = status.EventQueue
EventSubscription = status.EventSubscription
WGIfaceStatus = status.WGIfaceStatus
Listener = status.Listener
EventListener = status.EventListener
)
const (
StatusIdle = status.StatusIdle
StatusConnecting = status.StatusConnecting
StatusConnected = status.StatusConnected
)
var (
NewRecorder = status.NewRecorder
)
+4 -3
View File
@@ -9,6 +9,7 @@ import (
log "github.com/sirupsen/logrus" log "github.com/sirupsen/logrus"
"github.com/netbirdio/netbird/client/iface/configurer" "github.com/netbirdio/netbird/client/iface/configurer"
"github.com/netbirdio/netbird/client/internal/peer/status"
) )
type MocWgIface struct { type MocWgIface struct {
@@ -56,7 +57,7 @@ func TestWGWatcher_CheckSuccessCallback(t *testing.T) {
// platforms with coarse clock resolution (Windows), where two time.Now() calls // platforms with coarse clock resolution (Windows), where two time.Now() calls
// microseconds apart can return the same instant and read as a timed-out handshake. // microseconds apart can return the same instant and read as a timed-out handshake.
stats := &mockHandshakeStats{handshake: time.Now().Add(-time.Hour)} stats := &mockHandshakeStats{handshake: time.Now().Add(-time.Hour)}
watcher := NewWGWatcher(mlog, stats, "", newStateDump("peer", mlog, &Status{})) watcher := NewWGWatcher(mlog, stats, "", newStateDump("peer", mlog, &status.Recorder{}))
ctx, cancel := context.WithCancel(context.Background()) ctx, cancel := context.WithCancel(context.Background())
defer cancel() defer cancel()
@@ -104,7 +105,7 @@ func TestWGWatcher_EnableWgWatcher(t *testing.T) {
mlog := log.WithField("peer", "tet") mlog := log.WithField("peer", "tet")
mocWgIface := &MocWgIface{} mocWgIface := &MocWgIface{}
watcher := NewWGWatcher(mlog, mocWgIface, "", newStateDump("peer", mlog, &Status{})) watcher := NewWGWatcher(mlog, mocWgIface, "", newStateDump("peer", mlog, &status.Recorder{}))
ctx, cancel := context.WithCancel(context.Background()) ctx, cancel := context.WithCancel(context.Background())
defer cancel() defer cancel()
@@ -136,7 +137,7 @@ func TestWGWatcher_ReEnable(t *testing.T) {
mlog := log.WithField("peer", "tet") mlog := log.WithField("peer", "tet")
mocWgIface := &MocWgIface{} mocWgIface := &MocWgIface{}
watcher := NewWGWatcher(mlog, mocWgIface, "", newStateDump("peer", mlog, &Status{})) watcher := NewWGWatcher(mlog, mocWgIface, "", newStateDump("peer", mlog, &status.Recorder{}))
ctx, cancel := context.WithCancel(context.Background()) ctx, cancel := context.WithCancel(context.Background())
watcher.PrepareInitialHandshake() watcher.PrepareInitialHandshake()
+3 -2
View File
@@ -15,6 +15,7 @@ import (
"github.com/netbirdio/netbird/client/iface/udpmux" "github.com/netbirdio/netbird/client/iface/udpmux"
"github.com/netbirdio/netbird/client/internal/peer/conntype" "github.com/netbirdio/netbird/client/internal/peer/conntype"
icemaker "github.com/netbirdio/netbird/client/internal/peer/ice" icemaker "github.com/netbirdio/netbird/client/internal/peer/ice"
"github.com/netbirdio/netbird/client/internal/peer/status"
"github.com/netbirdio/netbird/client/internal/portforward" "github.com/netbirdio/netbird/client/internal/portforward"
"github.com/netbirdio/netbird/client/internal/stdnet" "github.com/netbirdio/netbird/client/internal/stdnet"
"github.com/netbirdio/netbird/route" "github.com/netbirdio/netbird/route"
@@ -39,7 +40,7 @@ type WorkerICE struct {
conn *Conn conn *Conn
signaler *Signaler signaler *Signaler
iFaceDiscover stdnet.ExternalIFaceDiscover iFaceDiscover stdnet.ExternalIFaceDiscover
statusRecorder *Status statusRecorder *status.Recorder
hasRelayOnLocally bool hasRelayOnLocally bool
agent *icemaker.ThreadSafeAgent agent *icemaker.ThreadSafeAgent
@@ -65,7 +66,7 @@ type WorkerICE struct {
portForwardAttempted bool portForwardAttempted bool
} }
func NewWorkerICE(ctx context.Context, log *log.Entry, config ConnConfig, conn *Conn, signaler *Signaler, ifaceDiscover stdnet.ExternalIFaceDiscover, statusRecorder *Status, hasRelayOnLocally bool) (*WorkerICE, error) { func NewWorkerICE(ctx context.Context, log *log.Entry, config ConnConfig, conn *Conn, signaler *Signaler, ifaceDiscover stdnet.ExternalIFaceDiscover, statusRecorder *status.Recorder, hasRelayOnLocally bool) (*WorkerICE, error) {
sessionID, err := NewICESessionID() sessionID, err := NewICESessionID()
if err != nil { if err != nil {
return nil, err return nil, err