mirror of
https://github.com/netbirdio/netbird.git
synced 2026-10-09 15:09:08 +02:00
Merge branch 'fix/subscribe-status-coalesce' into test/gui-memory-leak-fix
This commit is contained in:
@@ -1,20 +1,27 @@
|
|||||||
package server
|
package server
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"context"
|
||||||
|
"time"
|
||||||
|
|
||||||
log "github.com/sirupsen/logrus"
|
log "github.com/sirupsen/logrus"
|
||||||
|
|
||||||
"github.com/netbirdio/netbird/client/proto"
|
"github.com/netbirdio/netbird/client/proto"
|
||||||
)
|
)
|
||||||
|
|
||||||
|
const statusCoalesceWindow = 200 * time.Millisecond
|
||||||
|
|
||||||
// SubscribeStatus pushes a fresh StatusResponse on every connection state
|
// SubscribeStatus pushes a fresh StatusResponse on every connection state
|
||||||
// change. The first message is the current snapshot, so a re-subscribing
|
// change. The first message is the current snapshot, so a re-subscribing
|
||||||
// client doesn't need to also call Status. Subsequent messages fire when
|
// client doesn't need to also call Status. Subsequent messages fire when
|
||||||
// the peer recorder reports any of: connected/disconnected/connecting,
|
// the peer recorder reports any of: connected/disconnected/connecting,
|
||||||
// management or signal flip, address change, or peers list change.
|
// management or signal flip, address change, or peers list change.
|
||||||
//
|
//
|
||||||
// The change channel coalesces bursts to a single tick. If the consumer
|
// Bursts are coalesced deterministically: the first tick after a quiet
|
||||||
// is slow the daemon drops extras (not blocks), and the next snapshot
|
// period is sent immediately, then a short window swallows the rest of the
|
||||||
// the consumer pulls already reflects everything.
|
// burst and a single trailing snapshot covers whatever arrived meanwhile.
|
||||||
|
// Every send is a full snapshot of the recorder's current state, so
|
||||||
|
// swallowed ticks lose no information.
|
||||||
func (s *Server) SubscribeStatus(req *proto.StatusRequest, stream proto.DaemonService_SubscribeStatusServer) error {
|
func (s *Server) SubscribeStatus(req *proto.StatusRequest, stream proto.DaemonService_SubscribeStatusServer) error {
|
||||||
subID, ch := s.statusRecorder.SubscribeToStateChanges()
|
subID, ch := s.statusRecorder.SubscribeToStateChanges()
|
||||||
defer func() {
|
defer func() {
|
||||||
@@ -37,12 +44,50 @@ func (s *Server) SubscribeStatus(req *proto.StatusRequest, stream proto.DaemonSe
|
|||||||
if err := s.sendStatusSnapshot(req, stream); err != nil {
|
if err := s.sendStatusSnapshot(req, stream); err != nil {
|
||||||
return err
|
return err
|
||||||
}
|
}
|
||||||
|
pending, open := collectStatusBurst(stream.Context(), ch)
|
||||||
|
if pending {
|
||||||
|
if err := s.sendStatusSnapshot(req, stream); err != nil {
|
||||||
|
return err
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if !open {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
case <-stream.Context().Done():
|
case <-stream.Context().Done():
|
||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// collectStatusBurst waits out the coalesce window, absorbing further ticks.
|
||||||
|
// pending reports whether any tick arrived; open is false when the channel
|
||||||
|
// closed or the stream context ended.
|
||||||
|
func collectStatusBurst(ctx context.Context, ch <-chan struct{}) (pending, open bool) {
|
||||||
|
timer := time.NewTimer(statusCoalesceWindow)
|
||||||
|
defer timer.Stop()
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case _, ok := <-ch:
|
||||||
|
if !ok {
|
||||||
|
return pending, false
|
||||||
|
}
|
||||||
|
pending = true
|
||||||
|
case <-timer.C:
|
||||||
|
select {
|
||||||
|
case _, ok := <-ch:
|
||||||
|
if !ok {
|
||||||
|
return pending, false
|
||||||
|
}
|
||||||
|
pending = true
|
||||||
|
default:
|
||||||
|
}
|
||||||
|
return pending, true
|
||||||
|
case <-ctx.Done():
|
||||||
|
return false, false
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func (s *Server) sendStatusSnapshot(req *proto.StatusRequest, stream proto.DaemonService_SubscribeStatusServer) error {
|
func (s *Server) sendStatusSnapshot(req *proto.StatusRequest, stream proto.DaemonService_SubscribeStatusServer) error {
|
||||||
resp, err := s.buildStatusResponse(stream.Context(), req)
|
resp, err := s.buildStatusResponse(stream.Context(), req)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
|
|||||||
Reference in New Issue
Block a user