MCPcopy Create free account
hub / github.com/RTradeLtd/Lens / Run

Method Run

engine/queue/queue.go:98–121  ·  view source on GitHub ↗

Run maintains the queue and executes flushes as necessary

()

Source from the content-addressed store, hash-verified

96
97// Run maintains the queue and executes flushes as necessary
98func (q *Queue) Run() {
99 q.smux.Lock()
100 q.stopped = false
101 q.smux.Unlock()
102 q.l.Infow("spinning up queue", "rate", q.rate)
103 var ticker = time.NewTicker(q.rate)
104 for {
105 select {
106 case <-ticker.C:
107 q.flushIfNeeded()
108
109 case item := <-q.pendingC:
110 q.pendingItems[q.pending] = item
111 q.pending++
112 q.flushIfNeeded()
113
114 case <-q.stopC:
115 q.l.Infow("stopping background job")
116 ticker.Stop()
117 q.stop()
118 return
119 }
120 }
121}
122
123// Close stops the queue runner and releases queue assets
124func (q *Queue) Close() {

Callers 1

TestQueue_QueueFunction · 0.45

Calls 2

flushIfNeededMethod · 0.95
stopMethod · 0.95

Tested by 1

TestQueue_QueueFunction · 0.36