From 12f2e69af2077e266675bf571f60b80acda0c0e5 Mon Sep 17 00:00:00 2001 From: Viktor Liu Date: Wed, 24 Jun 2026 13:13:34 +0200 Subject: [PATCH] Log signal stall, ICE pair selection, restart cadence, sync content, and receive backpressure to attribute the regression --- client/internal/engine.go | 25 +++++++++- client/internal/peer/worker_ice.go | 2 +- shared/signal/client/grpc.go | 79 +++++++++++++++++++++++++++++- shared/signal/client/worker.go | 7 +++ 4 files changed, 109 insertions(+), 4 deletions(-) diff --git a/client/internal/engine.go b/client/internal/engine.go index 452075da8..42897613e 100644 --- a/client/internal/engine.go +++ b/client/internal/engine.go @@ -14,6 +14,7 @@ import ( "sort" "strings" "sync" + "sync/atomic" "time" "github.com/hashicorp/go-multierror" @@ -88,6 +89,13 @@ var ErrResetConnection = fmt.Errorf("reset connection") var ErrEngineAlreadyStarted = errors.New("engine already started") +// engineRestartCount and engineLastRestart track client-restart cadence across +// engine recreations so a restart loop is distinguishable from rare restarts. +var ( + engineRestartCount atomic.Int64 + engineLastRestart atomic.Int64 +) + type EngineConfig struct { WgPort int WgIfaceName string @@ -910,6 +918,14 @@ func (e *Engine) handleSync(update *mgmProto.SyncResponse) error { return e.ctx.Err() } + if nm := update.GetNetworkMap(); nm != nil { + log.Infof("sync update: serial=%d remotePeers=%d offlinePeers=%d routes=%d firewallRules=%d checks=%d configPresent=%v remotePeersEmpty=%v", + nm.GetSerial(), len(nm.GetRemotePeers()), len(nm.GetOfflinePeers()), len(nm.GetRoutes()), + len(nm.GetFirewallRules()), len(update.GetChecks()), update.GetNetbirdConfig() != nil, nm.GetRemotePeersIsEmpty()) + } else { + log.Infof("sync update: config-only (no network map), configPresent=%v", update.GetNetbirdConfig() != nil) + } + if update.NetworkMap != nil && update.NetworkMap.PeerConfig != nil { e.handleAutoUpdateVersion(update.NetworkMap.PeerConfig.AutoUpdate) } @@ -2171,7 +2187,14 @@ func (e *Engine) triggerClientRestart() { return } - log.Info("restarting engine") + // Cadence survives engine recreation (package-level), so a restart loop shows + // as a fast-climbing count with a short gap, distinct from rare intentional restarts. + n := engineRestartCount.Add(1) + var sinceLast time.Duration + if prev := engineLastRestart.Swap(time.Now().UnixNano()); prev != 0 { + sinceLast = time.Since(time.Unix(0, prev)) + } + log.Infof("restarting engine (restart #%d, %s since previous)", n, sinceLast.Round(time.Second)) CtxGetState(e.ctx).Set(StatusConnecting) _ = CtxGetState(e.ctx).Wrap(ErrResetConnection) log.Infof("cancelling client context, engine will be recreated") diff --git a/client/internal/peer/worker_ice.go b/client/internal/peer/worker_ice.go index b1aa3e0f9..34f0ab93d 100644 --- a/client/internal/peer/worker_ice.go +++ b/client/internal/peer/worker_ice.go @@ -461,7 +461,7 @@ func (w *WorkerICE) createForwardedCandidate(srflxCandidate ice.Candidate, mappi } func (w *WorkerICE) onICESelectedCandidatePair(agent *icemaker.ThreadSafeAgent, c1, c2 ice.Candidate) { - w.log.Debugf("selected candidate pair [local <-> remote] -> [%s <-> %s], peer %s", c1.String(), c2.String(), + w.log.Infof("selected candidate pair [local <-> remote] -> [%s <-> %s], peer %s", c1.String(), c2.String(), w.config.Key) pairStat, ok := agent.GetSelectedCandidatePairStats() diff --git a/shared/signal/client/grpc.go b/shared/signal/client/grpc.go index 611ab0c45..8000e39e7 100644 --- a/shared/signal/client/grpc.go +++ b/shared/signal/client/grpc.go @@ -85,6 +85,16 @@ type GrpcClient struct { // receive backpressure as a dead stream: reconnecting cannot help, since the // new stream feeds the same worker, and only triggers a reconnect storm. receiveHandoffBlocked atomic.Bool + // lastDecrypt holds the Unix-nano timestamp of the last message the decryption + // worker pulled off its queue. Diagnostic only: it lets a stall log show + // whether the worker was draining (busy) or idle when the stream went silent. + lastDecrypt atomic.Int64 + // handoffWaitTotal, handoffWaitMax (nanos) and handoffWaitCount accumulate the + // time the receive loop spent blocked handing messages to the worker. This is + // time not spent reading the stream, so it quantifies receive backpressure. + handoffWaitTotal atomic.Int64 + handoffWaitMax atomic.Int64 + handoffWaitCount atomic.Int64 } // NewClient creates a new Signal client @@ -360,6 +370,8 @@ func (c *GrpcClient) SendToStream(msg *proto.EncryptedMessage) error { // decryptMessage decrypts the body of the msg using Wireguard private key and Remote peer's public key func (c *GrpcClient) decryptMessage(msg *proto.EncryptedMessage) (*proto.Message, error) { + c.lastDecrypt.Store(time.Now().UnixNano()) + remoteKey, err := wgtypes.ParseKey(msg.GetKey()) if err != nil { return nil, err @@ -446,6 +458,12 @@ func (c *GrpcClient) idleSinceReceive() time.Duration { return time.Since(time.Unix(0, c.lastReceived.Load())) } +// idleSinceDecrypt returns how long since the worker last pulled a message. +// Diagnostic only: distinguishes a busy/wedged worker from an idle one. +func (c *GrpcClient) idleSinceDecrypt() time.Duration { + return time.Since(time.Unix(0, c.lastDecrypt.Load())) +} + // receiveAlive reports whether the receive stream shows liveness: it delivered a // frame within the inactivity threshold, or the receive loop is currently parked // handing a message to a busy decryption worker. In the latter case the loop has @@ -467,18 +485,55 @@ func (c *GrpcClient) watchReceiveStream(ctx context.Context, cancelStream contex defer ticker.Stop() var probeSentAt time.Time + var holdLogged bool + var statTicks int + var lastStatTotal int64 for { select { case <-ctx.Done(): return case <-ticker.C: + // Periodic backpressure summary so time lost to the worker handoff is + // visible even when no stall fires. Emitted ~once a minute and only + // when the wait grew, to stay quiet on a healthy stream. + if statTicks++; statTicks >= int(time.Minute/receiveWatchdogInterval) { + statTicks = 0 + if total, max, count := c.handoffWaitStats(); int64(total) > lastStatTotal { + log.Infof("signal receive backpressure: handoffWaitTotal=%s (+%s last min) handoffWaitMax=%s handoffMsgs=%d", + total.Round(time.Second), (total - time.Duration(lastStatTotal)).Round(time.Millisecond), + max.Round(time.Millisecond), count) + lastStatTotal = int64(total) + } + } + if c.receiveAlive() { + // Attribute the case that matters in the field: silent past the + // threshold but held because the receive loop is parked on the + // worker handoff (backpressure), not a dead stream. Log once per + // hold episode so a persistent worker stall is visible at info. + if c.idleSinceReceive() >= receiveInactivityThreshold && c.receiveHandoffBlocked.Load() { + if !holdLogged { + total, max, count := c.handoffWaitStats() + log.Infof("signal receive idle %s, loop blocked on worker handoff (idleDecrypt=%s queueDepth=%d connState=%s handoffWaitTotal=%s handoffWaitMax=%s handoffMsgs=%d); holding stream", + c.idleSinceReceive().Round(time.Second), c.idleSinceDecrypt().Round(time.Second), + c.decryptionWorker.QueueLen(), c.signalConn.GetState(), + total.Round(time.Second), max.Round(time.Millisecond), count) + holdLogged = true + } + } else { + holdLogged = false + } probeSentAt = time.Time{} continue } + holdLogged = false if !probeSentAt.IsZero() && time.Since(probeSentAt) >= receiveProbeTimeout { - log.Warnf("signal receive stream stalled: no messages for %s and probe did not return, reconnecting", c.idleSinceReceive().Round(time.Second)) + total, max, count := c.handoffWaitStats() + log.Warnf("signal receive stream stalled, reconnecting: idleRecv=%s idleDecrypt=%s handoffBlocked=%v queueDepth=%d connState=%s handoffWaitTotal=%s handoffWaitMax=%s handoffMsgs=%d probe did not return", + c.idleSinceReceive().Round(time.Second), c.idleSinceDecrypt().Round(time.Second), + c.receiveHandoffBlocked.Load(), c.decryptionWorker.QueueLen(), c.signalConn.GetState(), + total.Round(time.Second), max.Round(time.Millisecond), count) c.receiveStalled.Store(true) c.notifyDisconnected(errReceiveStreamStalled) cancelStream() @@ -536,15 +591,35 @@ func (c *GrpcClient) receive(stream proto.SignalExchange_ConnectStreamClient) er // The handoff blocks while the worker is busy, which parks this loop and // stops Recv. Flag it so the watchdog does not read the resulting silence - // as a dead stream. + // as a dead stream, and account the wait as receive backpressure. + handoffStart := time.Now() c.receiveHandoffBlocked.Store(true) if err := c.decryptionWorker.AddMsg(c.ctx, msg); err != nil { log.Errorf("failed to add message to decryption worker: %v", err) } c.receiveHandoffBlocked.Store(false) + c.recordHandoffWait(time.Since(handoffStart)) } } +// recordHandoffWait accumulates the time the receive loop was blocked handing a +// message to the worker. +func (c *GrpcClient) recordHandoffWait(d time.Duration) { + c.handoffWaitTotal.Add(int64(d)) + c.handoffWaitCount.Add(1) + for { + cur := c.handoffWaitMax.Load() + if int64(d) <= cur || c.handoffWaitMax.CompareAndSwap(cur, int64(d)) { + break + } + } +} + +// handoffWaitStats returns cumulative receive-loop handoff backpressure. +func (c *GrpcClient) handoffWaitStats() (total, max time.Duration, count int64) { + return time.Duration(c.handoffWaitTotal.Load()), time.Duration(c.handoffWaitMax.Load()), c.handoffWaitCount.Load() +} + func (c *GrpcClient) startEncryptionWorker(handler func(msg *proto.Message) error) { if c.decryptionWorker != nil { return diff --git a/shared/signal/client/worker.go b/shared/signal/client/worker.go index c724319b7..bafd38f39 100644 --- a/shared/signal/client/worker.go +++ b/shared/signal/client/worker.go @@ -32,6 +32,13 @@ func (w *Worker) AddMsg(ctx context.Context, msg *proto.EncryptedMessage) error return nil } +// QueueLen returns the number of messages buffered for decryption. Diagnostic +// only: a non-empty queue while the receive stream is silent indicates the +// receive loop is parked on the handoff rather than the stream being dead. +func (w *Worker) QueueLen() int { + return len(w.encryptedMsgPool) +} + func (w *Worker) Work(ctx context.Context) { for { select {