mirror of
https://github.com/netbirdio/netbird.git
synced 2026-09-27 00:59:07 +02:00
fixed a couple of issues flagged by coderabbit
Signed-off-by: Dmitri Dolguikh <dmitri.external@netbird.io>
This commit is contained in:
@@ -223,12 +223,12 @@ func (m *Manager) startSender(ctx context.Context, flowConfigInterval time.Durat
|
|||||||
collectedEvents := m.logger.ResetAggregationWindow()
|
collectedEvents := m.logger.ResetAggregationWindow()
|
||||||
events := collectedEvents.GetAggregatedEvents()
|
events := collectedEvents.GetAggregatedEvents()
|
||||||
for _, event := range events {
|
for _, event := range events {
|
||||||
|
m.eventsWithoutAcks.StoreEvent(event)
|
||||||
if err := m.send(event); err != nil {
|
if err := m.send(event); err != nil {
|
||||||
log.Errorf("failed to send flow event to server: %v", err)
|
log.Errorf("failed to send flow event to server: %v", err)
|
||||||
} else {
|
} else {
|
||||||
log.Tracef("sent flow event: %s", event.ID)
|
log.Tracef("sent flow event: %s", event.ID)
|
||||||
}
|
}
|
||||||
m.eventsWithoutAcks.StoreEvent(event)
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -271,6 +271,7 @@ func (m *Manager) startRetries(ctx context.Context, flowConfigInterval time.Dura
|
|||||||
case <-ctx.Done():
|
case <-ctx.Done():
|
||||||
return
|
return
|
||||||
case <-timer.C:
|
case <-timer.C:
|
||||||
|
resetBackoff := true
|
||||||
for _, e := range m.eventsWithoutAcks.GetEvents() {
|
for _, e := range m.eventsWithoutAcks.GetEvents() {
|
||||||
if e.Timestamp.Add(time.Second).After(time.Now()) {
|
if e.Timestamp.Add(time.Second).After(time.Now()) {
|
||||||
// grace period on retries to avoid early retries
|
// grace period on retries to avoid early retries
|
||||||
@@ -278,12 +279,15 @@ func (m *Manager) startRetries(ctx context.Context, flowConfigInterval time.Dura
|
|||||||
continue
|
continue
|
||||||
}
|
}
|
||||||
if err := m.send(e); err != nil {
|
if err := m.send(e); err != nil {
|
||||||
timer = time.NewTimer(retryBackoff.NextBackOff()) //nolint:staticcheck,wastedassign
|
timer = time.NewTimer(retryBackoff.NextBackOff())
|
||||||
|
resetBackoff = false
|
||||||
break
|
break
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
retryBackoff.Reset()
|
if resetBackoff { // use regular retry interval in absense of network errors
|
||||||
timer = time.NewTimer(m.retryInterval)
|
retryBackoff.Reset()
|
||||||
|
timer = time.NewTimer(m.retryInterval)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -70,7 +70,7 @@ func (am *AggregatingMemory) ResetAggregationWindow() types.FlowEventAggregator
|
|||||||
defer am.mux.Unlock()
|
defer am.mux.Unlock()
|
||||||
|
|
||||||
now := time.Now()
|
now := time.Now()
|
||||||
toret := AggregatingMemory{WindowStart: am.WindowStart, WindowEnd: now, Memory: Memory{events: am.events}}
|
toret := AggregatingMemory{WindowStart: am.WindowStart, WindowEnd: now, Memory: Memory{events: am.events}, rnd: v2.NewPCG(rand.Uint64(), rand.Uint64())}
|
||||||
|
|
||||||
am.events = make(map[uuid.UUID]*types.Event)
|
am.events = make(map[uuid.UUID]*types.Event)
|
||||||
am.WindowStart = now
|
am.WindowStart = now
|
||||||
|
|||||||
@@ -109,7 +109,7 @@ func (c *GRPCClient) Close() error {
|
|||||||
func (c *GRPCClient) Send(event *proto.FlowEvent) error {
|
func (c *GRPCClient) Send(event *proto.FlowEvent) error {
|
||||||
c.mu.Lock()
|
c.mu.Lock()
|
||||||
stream := c.stream
|
stream := c.stream
|
||||||
c.mu.Unlock()
|
defer c.mu.Unlock() // stream.Send() is not safe to call concurrently from multiple goroutines
|
||||||
|
|
||||||
if stream == nil {
|
if stream == nil {
|
||||||
return errors.New("stream not initialized")
|
return errors.New("stream not initialized")
|
||||||
|
|||||||
Reference in New Issue
Block a user