From 1f1413ec6a93b9f6469d93bcabee8e849d8f6cc5 Mon Sep 17 00:00:00 2001 From: Dmitri Dolguikh Date: Mon, 22 Jun 2026 18:34:08 +0200 Subject: [PATCH] responded to feedback + small fixes Signed-off-by: Dmitri Dolguikh --- client/internal/netflow/manager.go | 23 +++++++++++++---------- client/internal/netflow/store/memory.go | 14 +++++++++----- flow/proto/generate.sh | 1 - 3 files changed, 22 insertions(+), 16 deletions(-) diff --git a/client/internal/netflow/manager.go b/client/internal/netflow/manager.go index b1d0afecb..ba07c8522 100644 --- a/client/internal/netflow/manager.go +++ b/client/internal/netflow/manager.go @@ -72,6 +72,7 @@ func (m *Manager) needsNewClient(previous *nftypes.FlowConfig) bool { } // enableFlow starts components for flow tracking +// must be called under m.mux lock func (m *Manager) enableFlow(previous *nftypes.FlowConfig) error { // first make sender ready so events don't pile up if m.needsNewClient(previous) { @@ -91,6 +92,7 @@ func (m *Manager) enableFlow(previous *nftypes.FlowConfig) error { return nil } +// must be called under m.mux lock func (m *Manager) resetClient() error { if m.receiverClient != nil { if err := m.receiverClient.Close(); err != nil { @@ -114,17 +116,18 @@ func (m *Manager) resetClient() error { m.cancel = cancel m.shutdownWg.Add(3) + flowConfigInterval := m.flowConfig.Interval go func() { defer m.shutdownWg.Done() - m.receiveACKs(ctx, flowClient) + m.receiveACKs(ctx, flowClient, flowConfigInterval) }() go func() { defer m.shutdownWg.Done() - m.startSender(ctx) + m.startSender(ctx, flowConfigInterval) }() go func() { defer m.shutdownWg.Done() - m.startRetries(ctx) + m.startRetries(ctx, flowConfigInterval) }() return nil @@ -208,8 +211,8 @@ func (m *Manager) GetLogger() nftypes.FlowLogger { return m.logger } -func (m *Manager) startSender(ctx context.Context) { - ticker := time.NewTicker(m.flowConfig.Interval) +func (m *Manager) startSender(ctx context.Context, flowConfigInterval time.Duration) { + ticker := time.NewTicker(flowConfigInterval) defer ticker.Stop() for { @@ -231,8 +234,8 @@ func (m *Manager) startSender(ctx context.Context) { } } -func (m *Manager) receiveACKs(ctx context.Context, client *client.GRPCClient) { - err := client.Receive(ctx, m.flowConfig.Interval, func(ack *proto.FlowEventAck) error { +func (m *Manager) receiveACKs(ctx context.Context, client *client.GRPCClient, flowConfigInterval time.Duration) { + err := client.Receive(ctx, flowConfigInterval, func(ack *proto.FlowEventAck) error { id, err := uuid.FromBytes(ack.EventId) if err != nil { log.Warnf("failed to convert ack event id to uuid: %v", err) @@ -250,13 +253,13 @@ func (m *Manager) receiveACKs(ctx context.Context, client *client.GRPCClient) { // We effectively never drop events (see MaxInterval), which makes eventsWithoutAcks unbounded. // We may want to limit the max size of the store, and start dropping oldest events when the threshold is reached. -func (m *Manager) startRetries(ctx context.Context) { +func (m *Manager) startRetries(ctx context.Context, flowConfigInterval time.Duration) { timer := time.NewTimer(m.retryInterval) retryBackoff := backoff.WithContext(&backoff.ExponentialBackOff{ InitialInterval: 1 * time.Second, RandomizationFactor: 0.5, Multiplier: 1.7, - MaxInterval: m.flowConfig.Interval / 2, + MaxInterval: flowConfigInterval / 2, MaxElapsedTime: 3 * 30 * 24 * time.Hour, // 3 months Stop: backoff.Stop, Clock: backoff.SystemClock, @@ -280,7 +283,7 @@ func (m *Manager) startRetries(ctx context.Context) { } } retryBackoff.Reset() - timer = time.NewTimer(time.Second) + timer = time.NewTimer(m.retryInterval) } } } diff --git a/client/internal/netflow/store/memory.go b/client/internal/netflow/store/memory.go index 87a2f250b..aba2c6e48 100644 --- a/client/internal/netflow/store/memory.go +++ b/client/internal/netflow/store/memory.go @@ -2,6 +2,8 @@ package store import ( "maps" + "math/rand" + v2 "math/rand/v2" "net/netip" "slices" "sync" @@ -26,6 +28,7 @@ type AggregatingMemory struct { Memory WindowStart time.Time WindowEnd time.Time + rnd *v2.PCG } func (m *Memory) StoreEvent(event *types.Event) { @@ -59,17 +62,18 @@ func (m *Memory) DeleteEvents(ids []uuid.UUID) { } func NewAggregatingMemoryStore() *AggregatingMemory { - return &AggregatingMemory{WindowStart: time.Now(), Memory: Memory{events: make(map[uuid.UUID]*types.Event)}} + return &AggregatingMemory{WindowStart: time.Now(), Memory: Memory{events: make(map[uuid.UUID]*types.Event)}, rnd: v2.NewPCG(rand.Uint64(), rand.Uint64())} } func (am *AggregatingMemory) ResetAggregationWindow() types.FlowEventAggregator { am.mux.Lock() defer am.mux.Unlock() - toret := AggregatingMemory{WindowStart: am.WindowStart, WindowEnd: time.Now(), Memory: Memory{events: am.events}} + now := time.Now() + toret := AggregatingMemory{WindowStart: am.WindowStart, WindowEnd: now, Memory: Memory{events: am.events}} am.events = make(map[uuid.UUID]*types.Event) - am.WindowStart = time.Now() + am.WindowStart = now return &toret } @@ -81,7 +85,7 @@ type aggregationKey struct { direction int protocol uint8 icmpType uint8 - unique int64 // used to prevent aggregation on non icmp/udp/tcp events + unique uint64 // used to prevent aggregation on non icmp/udp/tcp events } func (am *AggregatingMemory) GetAggregatedEvents() []*types.Event { @@ -111,7 +115,7 @@ func (am *AggregatingMemory) GetAggregatedEvents() []*types.Event { event.WindowEnd = am.WindowEnd if event.Protocol != types.ICMP && event.Protocol != types.ICMPv6 && event.Protocol != types.UDP && event.Protocol != types.TCP { - lookupKey.unique = time.Now().UnixNano() // to make the lookup key unique so we don't aggregate on it + lookupKey.unique = am.rnd.Uint64() // to make the lookup key unique so we don't aggregate on it } aggregated[lookupKey] = event diff --git a/flow/proto/generate.sh b/flow/proto/generate.sh index fe447fd7a..6bbf78e61 100755 --- a/flow/proto/generate.sh +++ b/flow/proto/generate.sh @@ -10,7 +10,6 @@ fi old_pwd=$(pwd) script_path=$(dirname $(realpath "$0")) -echo "$script_path" cd "$script_path" go install google.golang.org/protobuf/cmd/protoc-gen-go@v1.26 go install google.golang.org/grpc/cmd/protoc-gen-go-grpc@v1.1