diff --git a/client/internal/peer/status.go b/client/internal/peer/status.go index 9b5fe16ee..00dd78b5f 100644 --- a/client/internal/peer/status.go +++ b/client/internal/peer/status.go @@ -228,6 +228,8 @@ type Status struct { routeIDLookup routeIDLookup wgIface WGIfaceStatus + + profile *StatusProfile } // NewRecorder returns a new Status instance @@ -242,16 +244,23 @@ func NewRecorder(mgmAddress string) *Status { notifier: newNotifier(), mgmAddress: mgmAddress, resolvedDomainsStates: map[domain.Domain]ResolvedDomainInfo{}, + profile: NewStatusProfile(context.Background()), } } +func (d *Status) StartProfile(ctx context.Context) { + d.profile.Start(ctx) +} + func (d *Status) SetRelayMgr(manager *relayClient.Manager) { + d.profile.inc("SetRelayMgr") d.muxRelays.Lock() defer d.muxRelays.Unlock() d.relayMgr = manager } func (d *Status) SetIngressGwMgr(ingressGwMgr *ingressgw.Manager) { + d.profile.inc("SetIngressGwMgr") d.mux.Lock() defer d.mux.Unlock() d.ingressGwMgr = ingressGwMgr @@ -259,6 +268,7 @@ func (d *Status) SetIngressGwMgr(ingressGwMgr *ingressgw.Manager) { // ReplaceOfflinePeers replaces func (d *Status) ReplaceOfflinePeers(replacement []State) { + d.profile.inc("ReplaceOfflinePeers") d.mux.Lock() defer d.mux.Unlock() d.offlinePeers = make([]State, len(replacement)) @@ -270,6 +280,7 @@ func (d *Status) ReplaceOfflinePeers(replacement []State) { // AddPeer adds peer to Daemon status map func (d *Status) AddPeer(peerPubKey string, fqdn string, ip string, ipv6 string) error { + d.profile.inc("AddPeer") d.mux.Lock() defer d.mux.Unlock() @@ -297,6 +308,7 @@ func (d *Status) AddPeer(peerPubKey string, fqdn string, ip string, ipv6 string) // GetPeer adds peer to Daemon status map func (d *Status) GetPeer(peerPubKey string) (State, error) { + d.profile.inc("GetPeer") d.mux.RLock() defer d.mux.RUnlock() @@ -308,6 +320,7 @@ func (d *Status) GetPeer(peerPubKey string) (State, error) { } func (d *Status) PeerByIP(ip string) (string, bool) { + d.profile.inc("PeerByIP") d.mux.RLock() defer d.mux.RUnlock() @@ -325,6 +338,7 @@ func (d *Status) PeerByIP(ip string) (string, bool) { // active peers are matched; peers moved into the offline slice by // ReplaceOfflinePeers are intentionally treated as unknown. func (d *Status) PeerStateByIP(ip string) (State, bool) { + d.profile.inc("PeerStateByIP") if ip == "" { return State{}, false } @@ -343,6 +357,7 @@ func (d *Status) PeerStateByIP(ip string) (State, bool) { // RemovePeer removes peer from Daemon status map func (d *Status) RemovePeer(peerPubKey string) error { + d.profile.inc("RemovePeer") d.mux.Lock() defer d.mux.Unlock() @@ -364,6 +379,7 @@ func (d *Status) RemovePeer(peerPubKey string) error { // UpdatePeerState updates peer status func (d *Status) UpdatePeerState(receivedState State) error { + d.profile.inc("UpdatePeerState") d.mux.Lock() peerState, ok := d.peers[receivedState.PubKey] @@ -406,6 +422,7 @@ func (d *Status) UpdatePeerState(receivedState State) error { } func (d *Status) AddPeerStateRoute(peer string, route string, resourceId route.ResID) error { + d.profile.inc("AddPeerStateRoute") d.mux.Lock() peerState, ok := d.peers[peer] @@ -431,6 +448,7 @@ func (d *Status) AddPeerStateRoute(peer string, route string, resourceId route.R } func (d *Status) RemovePeerStateRoute(peer string, route string) error { + d.profile.inc("RemovePeerStateRoute") d.mux.Lock() peerState, ok := d.peers[peer] @@ -461,11 +479,13 @@ func (d *Status) CheckRoutes(ip netip.Addr) ([]byte, bool) { if d == nil { return nil, false } + d.profile.inc("CheckRoutes") resId, isExitNode := d.routeIDLookup.Lookup(ip) return []byte(resId), isExitNode } func (d *Status) UpdatePeerICEState(receivedState State) error { + d.profile.inc("UpdatePeerICEState") d.mux.Lock() peerState, ok := d.peers[receivedState.PubKey] @@ -505,6 +525,7 @@ func (d *Status) UpdatePeerICEState(receivedState State) error { } func (d *Status) UpdatePeerRelayedState(receivedState State) error { + d.profile.inc("UpdatePeerRelayedState") d.mux.Lock() peerState, ok := d.peers[receivedState.PubKey] @@ -541,6 +562,7 @@ func (d *Status) UpdatePeerRelayedState(receivedState State) error { } func (d *Status) UpdatePeerRelayedStateToDisconnected(receivedState State) error { + d.profile.inc("UpdatePeerRelayedStateToDisconnected") d.mux.Lock() peerState, ok := d.peers[receivedState.PubKey] @@ -576,6 +598,7 @@ func (d *Status) UpdatePeerRelayedStateToDisconnected(receivedState State) error } func (d *Status) UpdatePeerICEStateToDisconnected(receivedState State) error { + d.profile.inc("UpdatePeerICEStateToDisconnected") d.mux.Lock() peerState, ok := d.peers[receivedState.PubKey] @@ -615,6 +638,7 @@ func (d *Status) UpdatePeerICEStateToDisconnected(receivedState State) error { // UpdateWireGuardPeerState updates the WireGuard bits of the peer state func (d *Status) UpdateWireGuardPeerState(pubKey string, wgStats configurer.WGStats) error { + d.profile.inc("UpdateWireGuardPeerState") d.mux.Lock() defer d.mux.Unlock() @@ -642,6 +666,7 @@ func hasConnStatusChanged(oldStatus, newStatus ConnStatus) bool { // UpdatePeerFQDN update peer's state fqdn only func (d *Status) UpdatePeerFQDN(peerPubKey, fqdn string) error { + d.profile.inc("UpdatePeerFQDN") d.mux.Lock() defer d.mux.Unlock() @@ -658,6 +683,7 @@ func (d *Status) UpdatePeerFQDN(peerPubKey, fqdn string) error { // UpdatePeerSSHHostKey updates peer's SSH host key func (d *Status) UpdatePeerSSHHostKey(peerPubKey string, sshHostKey []byte) error { + d.profile.inc("UpdatePeerSSHHostKey") d.mux.Lock() defer d.mux.Unlock() @@ -674,6 +700,7 @@ func (d *Status) UpdatePeerSSHHostKey(peerPubKey string, sshHostKey []byte) erro // FinishPeerListModifications this event invoke the notification func (d *Status) FinishPeerListModifications() { + d.profile.inc("FinishPeerListModifications") d.mux.Lock() if !d.peerListChangedForNotification { @@ -706,6 +733,7 @@ func (d *Status) FinishPeerListModifications() { } func (d *Status) SubscribeToPeerStateChanges(ctx context.Context, peerID string) *StatusChangeSubscription { + d.profile.inc("SubscribeToPeerStateChanges") d.mux.Lock() defer d.mux.Unlock() @@ -719,6 +747,7 @@ func (d *Status) SubscribeToPeerStateChanges(ctx context.Context, peerID string) } func (d *Status) UnsubscribePeerStateChanges(subscription *StatusChangeSubscription) { + d.profile.inc("UnsubscribePeerStateChanges") d.mux.Lock() defer d.mux.Unlock() @@ -744,6 +773,7 @@ func (d *Status) UnsubscribePeerStateChanges(subscription *StatusChangeSubscript // GetLocalPeerState returns the local peer state func (d *Status) GetLocalPeerState() LocalPeerState { + d.profile.inc("GetLocalPeerState") d.mux.RLock() defer d.mux.RUnlock() return d.localPeer.Clone() @@ -751,6 +781,7 @@ func (d *Status) GetLocalPeerState() LocalPeerState { // UpdateLocalPeerState updates local peer status func (d *Status) UpdateLocalPeerState(localPeerState LocalPeerState) { + d.profile.inc("UpdateLocalPeerState") d.mux.Lock() d.localPeer = localPeerState fqdn := d.localPeer.FQDN @@ -765,6 +796,7 @@ func (d *Status) UpdateLocalPeerState(localPeerState LocalPeerState) { // AddLocalPeerStateRoute adds a route to the local peer state func (d *Status) AddLocalPeerStateRoute(route string, resourceId route.ResID) { + d.profile.inc("AddLocalPeerStateRoute") d.mux.Lock() defer d.mux.Unlock() @@ -782,6 +814,7 @@ func (d *Status) AddLocalPeerStateRoute(route string, resourceId route.ResID) { // RemoveLocalPeerStateRoute removes a route from the local peer state func (d *Status) RemoveLocalPeerStateRoute(route string) { + d.profile.inc("RemoveLocalPeerStateRoute") d.mux.Lock() defer d.mux.Unlock() @@ -795,6 +828,7 @@ func (d *Status) RemoveLocalPeerStateRoute(route string) { // AddResolvedIPLookupEntry adds a resolved IP lookup entry func (d *Status) AddResolvedIPLookupEntry(prefix netip.Prefix, resourceId route.ResID) { + d.profile.inc("AddResolvedIPLookupEntry") d.mux.Lock() defer d.mux.Unlock() @@ -803,6 +837,7 @@ func (d *Status) AddResolvedIPLookupEntry(prefix netip.Prefix, resourceId route. // RemoveResolvedIPLookupEntry removes a resolved IP lookup entry func (d *Status) RemoveResolvedIPLookupEntry(route string) { + d.profile.inc("RemoveResolvedIPLookupEntry") d.mux.Lock() defer d.mux.Unlock() @@ -814,6 +849,7 @@ func (d *Status) RemoveResolvedIPLookupEntry(route string) { // CleanLocalPeerStateRoutes cleans all routes from the local peer state func (d *Status) CleanLocalPeerStateRoutes() { + d.profile.inc("CleanLocalPeerStateRoutes") d.mux.Lock() defer d.mux.Unlock() @@ -822,6 +858,7 @@ func (d *Status) CleanLocalPeerStateRoutes() { // CleanLocalPeerState cleans local peer status func (d *Status) CleanLocalPeerState() { + d.profile.inc("CleanLocalPeerState") d.mux.Lock() d.localPeer = LocalPeerState{} fqdn := d.localPeer.FQDN @@ -833,6 +870,7 @@ func (d *Status) CleanLocalPeerState() { // MarkManagementDisconnected sets ManagementState to disconnected func (d *Status) MarkManagementDisconnected(err error) { + d.profile.inc("MarkManagementDisconnected") d.mux.Lock() d.managementState = false d.managementError = err @@ -845,6 +883,7 @@ func (d *Status) MarkManagementDisconnected(err error) { // MarkManagementConnected sets ManagementState to connected func (d *Status) MarkManagementConnected() { + d.profile.inc("MarkManagementConnected") d.mux.Lock() d.managementState = true d.managementError = nil @@ -857,6 +896,7 @@ func (d *Status) MarkManagementConnected() { // UpdateSignalAddress update the address of the signal server func (d *Status) UpdateSignalAddress(signalURL string) { + d.profile.inc("UpdateSignalAddress") d.mux.Lock() defer d.mux.Unlock() d.signalAddress = signalURL @@ -864,6 +904,7 @@ func (d *Status) UpdateSignalAddress(signalURL string) { // UpdateManagementAddress update the address of the management server func (d *Status) UpdateManagementAddress(mgmAddress string) { + d.profile.inc("UpdateManagementAddress") d.mux.Lock() defer d.mux.Unlock() d.mgmAddress = mgmAddress @@ -871,6 +912,7 @@ func (d *Status) UpdateManagementAddress(mgmAddress string) { // UpdateRosenpass update the Rosenpass configuration func (d *Status) UpdateRosenpass(rosenpassEnabled, rosenpassPermissive bool) { + d.profile.inc("UpdateRosenpass") d.mux.Lock() defer d.mux.Unlock() d.rosenpassPermissive = rosenpassPermissive @@ -878,6 +920,7 @@ func (d *Status) UpdateRosenpass(rosenpassEnabled, rosenpassPermissive bool) { } func (d *Status) UpdateLazyConnection(enabled bool) { + d.profile.inc("UpdateLazyConnection") d.mux.Lock() defer d.mux.Unlock() d.lazyConnectionEnabled = enabled @@ -885,6 +928,7 @@ func (d *Status) UpdateLazyConnection(enabled bool) { // MarkSignalDisconnected sets SignalState to disconnected func (d *Status) MarkSignalDisconnected(err error) { + d.profile.inc("MarkSignalDisconnected") d.mux.Lock() d.signalState = false d.signalError = err @@ -897,6 +941,7 @@ func (d *Status) MarkSignalDisconnected(err error) { // MarkSignalConnected sets SignalState to connected func (d *Status) MarkSignalConnected() { + d.profile.inc("MarkSignalConnected") d.mux.Lock() d.signalState = true d.signalError = nil @@ -908,18 +953,21 @@ func (d *Status) MarkSignalConnected() { } func (d *Status) UpdateRelayStates(relayResults []relay.ProbeResult) { + d.profile.inc("UpdateRelayStates") d.muxRelays.Lock() defer d.muxRelays.Unlock() d.relayStates = relayResults } func (d *Status) UpdateDNSStates(dnsStates []NSGroupState) { + d.profile.inc("UpdateDNSStates") d.mux.Lock() defer d.mux.Unlock() d.nsGroupStates = dnsStates } func (d *Status) UpdateResolvedDomainsStates(originalDomain domain.Domain, resolvedDomain domain.Domain, prefixes []netip.Prefix, resourceId route.ResID) { + d.profile.inc("UpdateResolvedDomainsStates") d.mux.Lock() defer d.mux.Unlock() @@ -935,6 +983,7 @@ func (d *Status) UpdateResolvedDomainsStates(originalDomain domain.Domain, resol } func (d *Status) DeleteResolvedDomainsStates(domain domain.Domain) { + d.profile.inc("DeleteResolvedDomainsStates") d.mux.Lock() defer d.mux.Unlock() @@ -951,6 +1000,7 @@ func (d *Status) DeleteResolvedDomainsStates(domain domain.Domain) { } func (d *Status) GetRosenpassState() RosenpassState { + d.profile.inc("GetRosenpassState") d.mux.RLock() defer d.mux.RUnlock() return RosenpassState{ @@ -960,12 +1010,14 @@ func (d *Status) GetRosenpassState() RosenpassState { } func (d *Status) GetLazyConnection() bool { + d.profile.inc("GetLazyConnection") d.mux.RLock() defer d.mux.RUnlock() return d.lazyConnectionEnabled } func (d *Status) GetManagementState() ManagementState { + d.profile.inc("GetManagementState") d.mux.RLock() defer d.mux.RUnlock() return ManagementState{ @@ -976,6 +1028,7 @@ func (d *Status) GetManagementState() ManagementState { } func (d *Status) UpdateLatency(pubKey string, latency time.Duration) error { + d.profile.inc("UpdateLatency") if latency <= 0 { return nil } @@ -993,6 +1046,7 @@ func (d *Status) UpdateLatency(pubKey string, latency time.Duration) error { // IsLoginRequired determines if a peer's login has expired. func (d *Status) IsLoginRequired() bool { + d.profile.inc("IsLoginRequired") d.mux.RLock() defer d.mux.RUnlock() @@ -1009,6 +1063,7 @@ func (d *Status) IsLoginRequired() bool { } func (d *Status) GetSignalState() SignalState { + d.profile.inc("GetSignalState") d.mux.RLock() defer d.mux.RUnlock() return SignalState{ @@ -1020,6 +1075,7 @@ func (d *Status) GetSignalState() SignalState { // GetRelayStates returns the stun/turn/permanent relay states func (d *Status) GetRelayStates() []relay.ProbeResult { + d.profile.inc("GetRelayStates") d.muxRelays.RLock() // debug lines @@ -1064,6 +1120,7 @@ func (d *Status) GetRelayStates() []relay.ProbeResult { } func (d *Status) ForwardingRules() []firewall.ForwardRule { + d.profile.inc("ForwardingRules") d.mux.RLock() // debug lines @@ -1079,6 +1136,7 @@ func (d *Status) ForwardingRules() []firewall.ForwardRule { } func (d *Status) GetDNSStates() []NSGroupState { + d.profile.inc("GetDNSStates") d.mux.RLock() defer d.mux.RUnlock() @@ -1087,6 +1145,7 @@ func (d *Status) GetDNSStates() []NSGroupState { } func (d *Status) GetResolvedDomainsStates() map[domain.Domain]ResolvedDomainInfo { + d.profile.inc("GetResolvedDomainsStates") d.mux.RLock() defer d.mux.RUnlock() return maps.Clone(d.resolvedDomainsStates) @@ -1094,6 +1153,7 @@ func (d *Status) GetResolvedDomainsStates() map[domain.Domain]ResolvedDomainInfo // GetFullStatus gets full status func (d *Status) GetFullStatus() FullStatus { + d.profile.inc("GetFullStatus") fullStatus := FullStatus{ ManagementState: d.GetManagementState(), SignalState: d.GetSignalState(), @@ -1120,26 +1180,31 @@ func (d *Status) GetFullStatus() FullStatus { // ClientStart will notify all listeners about the new service state func (d *Status) ClientStart() { + d.profile.inc("ClientStart") d.notifier.clientStart() } // ClientStop will notify all listeners about the new service state func (d *Status) ClientStop() { + d.profile.inc("ClientStop") d.notifier.clientStop() } // ClientTeardown will notify all listeners about the service is under teardown func (d *Status) ClientTeardown() { + d.profile.inc("ClientTeardown") d.notifier.clientTearDown() } // SetConnectionListener set a listener to the notifier func (d *Status) SetConnectionListener(listener Listener) { + d.profile.inc("SetConnectionListener") d.notifier.setListener(listener) } // RemoveConnectionListener remove the listener from the notifier func (d *Status) RemoveConnectionListener() { + d.profile.inc("RemoveConnectionListener") d.notifier.removeListener() } @@ -1211,6 +1276,7 @@ func (d *Status) PublishEvent( userMsg string, metadata map[string]string, ) { + d.profile.inc("PublishEvent") event := &proto.SystemEvent{ Id: uuid.New().String(), Severity: severity, @@ -1244,6 +1310,7 @@ func (d *Status) PublishEvent( // SubscribeToEvents returns a new event subscription func (d *Status) SubscribeToEvents() *EventSubscription { + d.profile.inc("SubscribeToEvents") d.eventMux.Lock() defer d.eventMux.Unlock() @@ -1259,6 +1326,7 @@ func (d *Status) SubscribeToEvents() *EventSubscription { // UnsubscribeFromEvents removes an event subscription func (d *Status) UnsubscribeFromEvents(sub *EventSubscription) { + d.profile.inc("UnsubscribeFromEvents") if sub == nil { return } @@ -1274,10 +1342,12 @@ func (d *Status) UnsubscribeFromEvents(sub *EventSubscription) { // GetEventHistory returns all events in the queue func (d *Status) GetEventHistory() []*proto.SystemEvent { + d.profile.inc("GetEventHistory") return d.eventQueue.GetAll() } func (d *Status) SetWgIface(wgInterface WGIfaceStatus) { + d.profile.inc("SetWgIface") d.mux.Lock() defer d.mux.Unlock() @@ -1285,6 +1355,7 @@ func (d *Status) SetWgIface(wgInterface WGIfaceStatus) { } func (d *Status) PeersStatus() (*configurer.Stats, error) { + d.profile.inc("PeersStatus") d.mux.RLock() // debug lines @@ -1303,6 +1374,7 @@ func (d *Status) PeersStatus() (*configurer.Stats, error) { // and updates the cached peer states. This ensures accurate handshake times and // transfer statistics in status reports without running full health probes. func (d *Status) RefreshWireGuardStats() error { + d.profile.inc("RefreshWireGuardStats") d.mux.Lock() // debug lines diff --git a/client/internal/peer/status_profile.go b/client/internal/peer/status_profile.go new file mode 100644 index 000000000..45faa7794 --- /dev/null +++ b/client/internal/peer/status_profile.go @@ -0,0 +1,97 @@ +package peer + +import ( + "context" + "sort" + "strconv" + "strings" + "sync" + "sync/atomic" + "time" + + log "github.com/sirupsen/logrus" +) + +type StatusProfile struct { + counts sync.Map +} + +func NewStatusProfile(ctx context.Context) *StatusProfile { + s := &StatusProfile{} + go s.Start(ctx) + return s +} + +func (s *StatusProfile) Start(ctx context.Context) { + ticker := time.NewTicker(10 * time.Second) + defer ticker.Stop() + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + s.logCounts() + } + } +} + +func (s *StatusProfile) inc(method string) { + if s == nil { + return + } + if v, ok := s.counts.Load(method); ok { + v.(*atomic.Int64).Add(1) + return + } + cnt := &atomic.Int64{} + actual, _ := s.counts.LoadOrStore(method, cnt) + actual.(*atomic.Int64).Add(1) +} + +func (s *StatusProfile) snapshot() map[string]int64 { + out := make(map[string]int64) + s.counts.Range(func(k, v any) bool { + out[k.(string)] = v.(*atomic.Int64).Load() + return true + }) + return out +} + +func (s *StatusProfile) logCounts() { + counts := s.snapshot() + if len(counts) == 0 { + log.Infof("status profile: no Status method calls so far") + return + } + + type kv struct { + method string + count int64 + } + sorted := make([]kv, 0, len(counts)) + var total int64 + for m, c := range counts { + if c == 0 { + continue + } + sorted = append(sorted, kv{m, c}) + total += c + } + sort.Slice(sorted, func(i, j int) bool { + if sorted[i].count != sorted[j].count { + return sorted[i].count > sorted[j].count + } + return sorted[i].method < sorted[j].method + }) + + var b strings.Builder + for i, e := range sorted { + if i > 0 { + b.WriteString(", ") + } + b.WriteString(e.method) + b.WriteByte('=') + b.WriteString(strconv.FormatInt(e.count, 10)) + } + log.Infof("status profile (cumulative total=%d): %s", total, b.String()) +} diff --git a/client/server/server.go b/client/server/server.go index 3f6dabc56..176af1461 100644 --- a/client/server/server.go +++ b/client/server/server.go @@ -136,6 +136,7 @@ func New(ctx context.Context, logFile string, configFile string, profilesDisable networksDisabled: networksDisabled, jwtCache: newJWTCache(), } + go s.statusRecorder.StartProfile(ctx) agent := &serverAgent{s} s.sleepHandler = sleephandler.New(agent) s.startSleepDetector()