@@ -38,6 +38,11 @@ const (
3838 poolContextEventPosition poolContext = "eventpos"
3939)
4040
41+ type spoolEntry struct {
42+ events []* event.Event
43+ size int
44+ }
45+
4146// Pool manages a list of receivers
4247type Pool struct {
4348 // Pipeline
@@ -53,6 +58,8 @@ type Pool struct {
5358 scheduler * scheduler.Scheduler
5459 connectionLock sync.RWMutex
5560 connectionStatus map [interface {}]* poolConnectionStatus
61+ spool []* spoolEntry
62+ spoolSize int64
5663
5764 apiConfig * admin.Config
5865 apiConnections api.Array
@@ -104,16 +111,15 @@ func (r *Pool) Init(cfg *config.Config) error {
104111
105112// Run starts listening
106113func (r * Pool ) Run () {
107- var spool [][]* event.Event
108114 var spoolChan chan <- []* event.Event
109115 eventChan := r .eventChan
110116 shutdownChan := r .shutdownChan
111117
112118ReceiverLoop:
113119 for {
114- var nextSpool [] * event. Event = nil
115- if len (spool ) != 0 {
116- nextSpool = spool [0 ]
120+ var nextSpool * spoolEntry = nil
121+ if len (r . spool ) != 0 {
122+ nextSpool = r . spool [0 ]
117123 }
118124
119125 select {
@@ -182,30 +188,52 @@ ReceiverLoop:
182188 // We replace because a reconnect on the same port could occur before we get around to handling the disconnection, and we're keyed by port
183189 r .apiConnections .ReplaceEntry (eventImpl .Remote (), connectionStatus )
184190 case transports.EventsEvent :
191+ size := calcSize (eventImpl )
185192 r .connectionLock .Lock ()
186193 connection := eventImpl .Context ().Value (transports .ContextConnection )
187194 receiver := eventImpl .Context ().Value (transports .ContextReceiver ).(transports.Receiver )
188195 connectionStatus := r .connectionStatus [connection ]
189- // Schedule partial ack if this is first set of events
190- if len (connectionStatus .progress ) == 0 {
191- r .scheduler .Set (connection , 5 * time .Second )
196+ if r .spoolSize + int64 (size ) > r .receivers [receiver ].config .MaxQueueSize {
197+ receiver .ShutdownConnectionRead (eventImpl .Context (), fmt .Errorf ("max queue size exceeded" ))
198+ r .connectionLock .Unlock ()
199+ break
192200 }
193- connectionStatus .progress = append (connectionStatus .progress , & poolEventProgress {event : eventImpl , sequence : 0 })
201+ var acker event.Acknowledger
202+ if receiver .SupportsAck () {
203+ if len (r .connectionStatus [connection ].progress )+ 1 > int (r .receivers [receiver ].config .MaxPendingPayloads ) {
204+ receiver .ShutdownConnectionRead (eventImpl .Context (), fmt .Errorf ("max pending payloads exceeded" ))
205+ r .connectionLock .Unlock ()
206+ break
207+ }
208+ // Schedule partial ack if this is first set of events
209+ if len (connectionStatus .progress ) == 0 {
210+ r .scheduler .Set (connection , 5 * time .Second )
211+ }
212+ connectionStatus .progress = append (connectionStatus .progress , & poolEventProgress {event : eventImpl , sequence : 0 })
213+ acker = r
214+ } else {
215+ // Reset idle timeout
216+ r .startIdleTimeout (eventImpl .Context (), receiver , connection )
217+ }
218+ connectionStatus .bytes += eventImpl .Size ()
194219 r .connectionLock .Unlock ()
195220 // Build the events with our acknowledger and submit the bundle
196221 var events = make ([]* event.Event , len (eventImpl .Events ()))
222+ var ctx context.Context
197223 for idx , item := range eventImpl .Events () {
198- ctx := context .WithValue (eventImpl .Context (), poolContextEventPosition , & poolEventPosition {nonce : eventImpl .Nonce (), sequence : uint32 (idx + 1 )})
199- item := event .NewEvent (ctx , r , item )
224+ if acker == nil {
225+ ctx = eventImpl .Context ()
226+ } else {
227+ ctx = context .WithValue (eventImpl .Context (), poolContextEventPosition , & poolEventPosition {nonce : eventImpl .Nonce (), sequence : uint32 (idx + 1 )})
228+ }
229+ item := event .NewEvent (ctx , acker , item )
200230 item .MustResolve ("@metadata[receiver]" , connectionStatus .metadataReceiver )
201231 events [idx ] = item
202232 }
203- spool = append (spool , events )
233+ spoolEntry := & spoolEntry {events , size }
234+ r .spool = append (r .spool , spoolEntry )
235+ r .spoolSize += int64 (spoolEntry .size )
204236 spoolChan = r .output
205- // Stop reading events if this client breached our limit
206- if len (r .connectionStatus [connection ].progress ) > int (r .receivers [receiver ].config .MaxPendingPayloads ) {
207- receiver .ShutdownConnectionRead (eventImpl .Context (), fmt .Errorf ("max pending payloads exceeded" ))
208- }
209237 case * transports.EndEvent :
210238 // Connection EOF
211239 r .connectionLock .Lock ()
@@ -263,10 +291,11 @@ ReceiverLoop:
263291 r .startIdleTimeout (eventImpl .Context (), receiver , connection )
264292 }
265293 }
266- case spoolChan <- nextSpool :
267- copy (spool , spool [1 :])
268- spool = spool [:len (spool )- 1 ]
269- if len (spool ) == 0 {
294+ case spoolChan <- nextSpool .events :
295+ copy (r .spool , r .spool [1 :])
296+ r .spool = r .spool [:len (r .spool )- 1 ]
297+ r .spoolSize -= int64 (nextSpool .size )
298+ if len (r .spool ) == 0 {
270299 spoolChan = nil
271300 }
272301 }
@@ -385,6 +414,7 @@ func (r *Pool) updateReceivers(newConfig *config.Config) {
385414 receiverApi := & api.KeyValue {}
386415 receiverApi .SetEntry ("listen" , api .String (listen ))
387416 receiverApi .SetEntry ("maxPendingPayloads" , api .Number (cfgEntry .MaxPendingPayloads ))
417+ receiverApi .SetEntry ("maxQueueSize" , api .Number (cfgEntry .MaxQueueSize ))
388418 r .apiListeners .AddEntry (listen , receiverApi )
389419 }
390420 }
@@ -414,3 +444,7 @@ func (r *Pool) shutdown() {
414444 receiver .Shutdown ()
415445 }
416446}
447+
448+ func calcSize (eventImpl transports.EventsEvent ) int {
449+ return eventImpl .Size ()
450+ }
0 commit comments