MCPcopy Create free account
hub / github.com/TalkingData/owl / processSingleQueue

Function processSingleQueue

controller/notify.go:52–71  ·  view source on GitHub ↗

TODO: fix when product delete, goroutine leak

(queue *EventPool)

Source from the content-addressed store, hash-verified

50
51//TODO: fix when product delete, goroutine leak
52func processSingleQueue(queue *EventPool) {
53 duration := time.Millisecond * 100
54 for {
55 if queue.len() > GlobalConfig.SEND_MAX {
56 duration = time.Microsecond * time.Duration(queue.len())
57 if duration.Seconds() > float64(GlobalConfig.MAX_INTERVAL_WAIT_TIME) {
58 duration = time.Second * time.Duration(GlobalConfig.MAX_INTERVAL_WAIT_TIME)
59 }
60 }
61 if !queue.mute {
62 event := queue.getQueueEvent()
63 if queue.mute {
64 queue.putQueueEvent(event)
65 } else {
66 go processSingleEvent(event)
67 }
68 }
69 time.Sleep(duration)
70 }
71}
72
73func processSingleEvent(event *QueueEvent) {
74 switch event.status {

Callers 2

refreshQueueMethod · 0.85

Calls 4

processSingleEventFunction · 0.85
lenMethod · 0.80
getQueueEventMethod · 0.80
putQueueEventMethod · 0.80

Tested by

no test coverage detected