11package containerwatcher
22
33import (
4+ "sync"
45 "time"
56
67 "github.com/kubescape/go-logger"
78 "github.com/kubescape/go-logger/helpers"
89 "github.com/kubescape/node-agent/pkg/utils"
9- "github.com/oleiade/lane/v2"
1010)
1111
1212type EventEntry struct {
@@ -19,14 +19,15 @@ type EventEntry struct {
1919
2020type OrderedEventQueue struct {
2121 maxBufferSize int
22- eventQueue * lane.PriorityQueue [EventEntry , int64 ]
22+ eventQueue []EventEntry
23+ mutex sync.Mutex
2324 fullQueueAlert chan struct {}
2425}
2526
2627func NewOrderedEventQueue (collectionInterval time.Duration , maxBufferSize int ) * OrderedEventQueue {
2728 return & OrderedEventQueue {
2829 maxBufferSize : maxBufferSize ,
29- eventQueue : lane . NewMinPriorityQueue [ EventEntry , int64 ]( ),
30+ eventQueue : make ([] EventEntry , 0 , 1024 ),
3031 fullQueueAlert : make (chan struct {}, 1 ),
3132 }
3233}
@@ -36,11 +37,15 @@ func (oeq *OrderedEventQueue) GetFullQueueAlertChannel() <-chan struct{} {
3637}
3738
3839func (oeq * OrderedEventQueue ) AddEventDirect (eventType utils.EventType , event utils.K8sEvent , containerID string , processID uint32 ) {
39- if int (oeq .eventQueue .Size ()) >= oeq .maxBufferSize {
40+ oeq .mutex .Lock ()
41+ if len (oeq .eventQueue ) >= oeq .maxBufferSize {
42+ queueSize := len (oeq .eventQueue )
43+ oeq .mutex .Unlock ()
44+
4045 logger .L ().Warning ("Ordered event queue - Event queue full, dropping event to prevent OOM" ,
4146 helpers .String ("eventType" , string (eventType )),
4247 helpers .String ("containerID" , containerID ),
43- helpers .Int ("queueSize" , int ( oeq . eventQueue . Size ()) ),
48+ helpers .Int ("queueSize" , queueSize ),
4449 helpers .Int ("maxBufferSize" , oeq .maxBufferSize ))
4550
4651 select {
@@ -62,34 +67,98 @@ func (oeq *OrderedEventQueue) AddEventDirect(eventType utils.EventType, event ut
6267 ProcessID : processID ,
6368 }
6469
65- priority := timestamp . UnixNano ( )
66- oeq .eventQueue . Push ( eventEntry , priority )
70+ oeq . pushHeap ( eventEntry )
71+ oeq .mutex . Unlock ( )
6772}
6873
6974func (oeq * OrderedEventQueue ) PopEvent () (EventEntry , bool ) {
70- if oeq .eventQueue .Empty () {
75+ oeq .mutex .Lock ()
76+ defer oeq .mutex .Unlock ()
77+
78+ if len (oeq .eventQueue ) == 0 {
7179 return EventEntry {}, false
7280 }
7381
74- event , _ , ok := oeq .eventQueue .Pop ()
75- return event , ok
82+ return oeq .popHeap (), true
83+ }
84+
85+ func (oeq * OrderedEventQueue ) PopBatch (maxCount int , dst []EventEntry ) []EventEntry {
86+ oeq .mutex .Lock ()
87+ defer oeq .mutex .Unlock ()
88+
89+ dst = dst [:0 ]
90+ for len (oeq .eventQueue ) > 0 && len (dst ) < maxCount {
91+ dst = append (dst , oeq .popHeap ())
92+ }
93+ return dst
7694}
7795
7896func (oeq * OrderedEventQueue ) PeekEvent () (EventEntry , bool ) {
79- if oeq .eventQueue .Empty () {
97+ oeq .mutex .Lock ()
98+ defer oeq .mutex .Unlock ()
99+
100+ if len (oeq .eventQueue ) == 0 {
80101 return EventEntry {}, false
81102 }
82103
83- event , _ , ok := oeq .eventQueue .Head ()
84- return event , ok
104+ return oeq .eventQueue [0 ], true
85105}
86106
87107// Size returns the number of events in the queue
88108func (oeq * OrderedEventQueue ) Size () int {
89- return int (oeq .eventQueue .Size ())
109+ oeq .mutex .Lock ()
110+ defer oeq .mutex .Unlock ()
111+ return len (oeq .eventQueue )
90112}
91113
92114// Empty returns whether the queue is empty
93115func (oeq * OrderedEventQueue ) Empty () bool {
94- return oeq .eventQueue .Empty ()
116+ oeq .mutex .Lock ()
117+ defer oeq .mutex .Unlock ()
118+ return len (oeq .eventQueue ) == 0
119+ }
120+
121+ func (oeq * OrderedEventQueue ) pushHeap (entry EventEntry ) {
122+ oeq .eventQueue = append (oeq .eventQueue , entry )
123+ oeq .up (len (oeq .eventQueue ) - 1 )
124+ }
125+
126+ func (oeq * OrderedEventQueue ) popHeap () EventEntry {
127+ n := len (oeq .eventQueue ) - 1
128+ oeq .eventQueue [0 ], oeq .eventQueue [n ] = oeq .eventQueue [n ], oeq .eventQueue [0 ]
129+ oeq .down (0 , n )
130+ x := oeq .eventQueue [n ]
131+ oeq .eventQueue [n ] = EventEntry {}
132+ oeq .eventQueue = oeq .eventQueue [:n ]
133+ return x
134+ }
135+
136+ func (oeq * OrderedEventQueue ) up (j int ) {
137+ for {
138+ i := (j - 1 ) / 2 // parent
139+ if i == j || ! oeq .eventQueue [j ].Timestamp .Before (oeq .eventQueue [i ].Timestamp ) {
140+ break
141+ }
142+ oeq .eventQueue [i ], oeq .eventQueue [j ] = oeq .eventQueue [j ], oeq .eventQueue [i ]
143+ j = i
144+ }
145+ }
146+
147+ func (oeq * OrderedEventQueue ) down (i0 , n int ) {
148+ i := i0
149+ for {
150+ j1 := 2 * i + 1
151+ if j1 >= n || j1 < 0 {
152+ break
153+ }
154+ j := j1
155+ if j2 := j1 + 1 ; j2 < n && oeq .eventQueue [j2 ].Timestamp .Before (oeq .eventQueue [j1 ].Timestamp ) {
156+ j = j2
157+ }
158+ if ! oeq .eventQueue [j ].Timestamp .Before (oeq .eventQueue [i ].Timestamp ) {
159+ break
160+ }
161+ oeq .eventQueue [i ], oeq .eventQueue [j ] = oeq .eventQueue [j ], oeq .eventQueue [i ]
162+ i = j
163+ }
95164}
0 commit comments