diff --git a/client/internal/netflow/logger/logger.go b/client/internal/netflow/logger/logger.go index 8f8e68784..deb38bc4d 100644 --- a/client/internal/netflow/logger/logger.go +++ b/client/internal/netflow/logger/logger.go @@ -27,7 +27,7 @@ type Logger struct { wgIfaceNetV6 netip.Prefix dnsCollection atomic.Bool exitNodeCollection atomic.Bool - Store types.Store + Store types.AggregatingStore } func New(statusRecorder *peer.Status, wgIfaceIPNet, wgIfaceIPNetV6 netip.Prefix) *Logger { @@ -35,7 +35,7 @@ func New(statusRecorder *peer.Status, wgIfaceIPNet, wgIfaceIPNetV6 netip.Prefix) statusRecorder: statusRecorder, wgIfaceNet: wgIfaceIPNet, wgIfaceNetV6: wgIfaceIPNetV6, - Store: store.NewMemoryStore(), + Store: store.NewAggregatingMemoryStore(), } } @@ -125,6 +125,10 @@ func (l *Logger) stop() { l.mux.Unlock() } +func (l *Logger) ResetAggregationWindow() types.FlowEventAggregator { + return l.Store.ResetAggregationWindow() +} + func (l *Logger) GetEvents() []*types.Event { return l.Store.GetEvents() } diff --git a/client/internal/netflow/manager.go b/client/internal/netflow/manager.go index eff083dbf..f6397952d 100644 --- a/client/internal/netflow/manager.go +++ b/client/internal/netflow/manager.go @@ -9,12 +9,14 @@ import ( "sync" "time" + "github.com/cenkalti/backoff/v4" "github.com/google/uuid" log "github.com/sirupsen/logrus" "google.golang.org/protobuf/types/known/timestamppb" "github.com/netbirdio/netbird/client/internal/netflow/conntrack" "github.com/netbirdio/netbird/client/internal/netflow/logger" + "github.com/netbirdio/netbird/client/internal/netflow/store" nftypes "github.com/netbirdio/netbird/client/internal/netflow/types" "github.com/netbirdio/netbird/client/internal/peer" "github.com/netbirdio/netbird/flow/client" @@ -23,14 +25,15 @@ import ( // Manager handles netflow tracking and logging type Manager struct { - mux sync.Mutex - shutdownWg sync.WaitGroup - logger nftypes.FlowLogger - flowConfig *nftypes.FlowConfig - conntrack nftypes.ConnTracker - receiverClient *client.GRPCClient - publicKey []byte - cancel context.CancelFunc + mux sync.Mutex + shutdownWg sync.WaitGroup + logger nftypes.FlowLogger + flowConfig *nftypes.FlowConfig + conntrack nftypes.ConnTracker + receiverClient *client.GRPCClient + eventsWithoutAcks nftypes.Store + publicKey []byte + cancel context.CancelFunc } // NewManager creates a new netflow manager @@ -48,9 +51,10 @@ func NewManager(iface nftypes.IFaceMapper, publicKey []byte, statusRecorder *pee } return &Manager{ - logger: flowLogger, - conntrack: ct, - publicKey: publicKey, + logger: flowLogger, + conntrack: ct, + publicKey: publicKey, + eventsWithoutAcks: store.NewMemoryStore(), } } @@ -107,7 +111,7 @@ func (m *Manager) resetClient() error { ctx, cancel := context.WithCancel(context.Background()) m.cancel = cancel - m.shutdownWg.Add(2) + m.shutdownWg.Add(3) go func() { defer m.shutdownWg.Done() m.receiveACKs(ctx, flowClient) @@ -116,6 +120,10 @@ func (m *Manager) resetClient() error { defer m.shutdownWg.Done() m.startSender(ctx) }() + go func() { + defer m.shutdownWg.Done() + m.startRetries(ctx) + }() return nil } @@ -207,13 +215,16 @@ func (m *Manager) startSender(ctx context.Context) { case <-ctx.Done(): return case <-ticker.C: - events := m.logger.GetEvents() + collectedEvents := m.logger.ResetAggregationWindow() + events := collectedEvents.GetAggregatedEvents() for _, event := range events { + // handle retries, grace period? if err := m.send(event); err != nil { log.Errorf("failed to send flow event to server: %v", err) - continue + } else { + log.Tracef("sent flow event: %s", event.ID) } - log.Tracef("sent flow event: %s", event.ID) + m.eventsWithoutAcks.StoreEvent(event) } } } @@ -227,7 +238,7 @@ func (m *Manager) receiveACKs(ctx context.Context, client *client.GRPCClient) { return nil } log.Tracef("received flow event ack: %s", id) - m.logger.DeleteEvents([]uuid.UUID{id}) + m.eventsWithoutAcks.DeleteEvents([]uuid.UUID{id}) return nil }) @@ -236,6 +247,36 @@ func (m *Manager) receiveACKs(ctx context.Context, client *client.GRPCClient) { } } +func (m *Manager) startRetries(ctx context.Context) { + ticker := time.NewTimer(time.Second) + retryBackoff := backoff.WithContext(&backoff.ExponentialBackOff{ + InitialInterval: 1 * time.Second, + RandomizationFactor: 0.5, + Multiplier: 1.7, + MaxInterval: m.flowConfig.Interval / 2, + MaxElapsedTime: 3 * 30 * 24 * time.Hour, // 3 months + Stop: backoff.Stop, + Clock: backoff.SystemClock, + }, ctx) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + for _, e := range m.eventsWithoutAcks.GetEvents() { + if err := m.send(e); err != nil { + ticker = time.NewTimer(retryBackoff.NextBackOff()) + break + } + } + retryBackoff.Reset() + ticker = time.NewTimer(time.Second) + } + } +} + func (m *Manager) send(event *nftypes.Event) error { m.mux.Lock() client := m.receiverClient diff --git a/client/internal/netflow/store/memory.go b/client/internal/netflow/store/memory.go index f9151726e..92e4578fc 100644 --- a/client/internal/netflow/store/memory.go +++ b/client/internal/netflow/store/memory.go @@ -57,11 +57,15 @@ func (m *Memory) DeleteEvents(ids []uuid.UUID) { } } -func (am *AggregatingMemory) StartAggregationWindow() *AggregatingMemory { +func NewAggregatingMemoryStore() *AggregatingMemory { + return &AggregatingMemory{Memory{events: make(map[uuid.UUID]*types.Event)}} +} + +func (am *AggregatingMemory) ResetAggregationWindow() types.FlowEventAggregator { am.mux.Lock() defer am.mux.Unlock() - toret := AggregatingMemory{Memory: Memory{events: am.Memory.events}} + toret := AggregatingMemory{Memory: Memory{events: am.events}} am.events = make(map[uuid.UUID]*types.Event) return &toret @@ -82,6 +86,7 @@ func (am *AggregatingMemory) GetAggregatedEvents() []*types.Event { if aggregatedEvent, ok := aggregated[lookupKey]; ok { switch aggregatedEvent.Protocol { case types.ICMP, types.ICMPv6, types.UDP, types.TCP: + // track the number of connections, duration?, open and close events? aggregatedEvent.RxBytes += v.RxBytes aggregatedEvent.RxPackets += v.RxPackets aggregatedEvent.TxBytes += v.TxBytes @@ -98,7 +103,7 @@ func (am *AggregatingMemory) GetAggregatedEvents() []*types.Event { case types.ICMP, types.ICMPv6, types.TCP, types.UDP: aggregated[lookupKey] = v default: - lookupKey.ts = time.Now().UnixNano() + lookupKey.ts = time.Now().UnixNano() // to make the lookup key unique so we don't aggregate on it aggregated[lookupKey] = v } } diff --git a/client/internal/netflow/types/types.go b/client/internal/netflow/types/types.go index 3f7d0d0ad..55c7494be 100644 --- a/client/internal/netflow/types/types.go +++ b/client/internal/netflow/types/types.go @@ -114,13 +114,15 @@ type FlowManager interface { GetLogger() FlowLogger } +type FlowEventAggregator interface { + ResetAggregationWindow() FlowEventAggregator + GetAggregatedEvents() []*Event +} + type FlowLogger interface { + ResetAggregationWindow() FlowEventAggregator // StoreEvent stores a flow event StoreEvent(flowEvent EventFields) - // GetEvents returns all stored events - GetEvents() []*Event - // DeleteEvents deletes events from the store - DeleteEvents([]uuid.UUID) // Close closes the logger Close() // Enable enables the flow logger receiver @@ -140,6 +142,11 @@ type Store interface { Close() } +type AggregatingStore interface { + FlowEventAggregator + Store +} + // ConnTracker defines the interface for connection tracking functionality type ConnTracker interface { // Start begins tracking connections by listening for conntrack events.