-
Notifications
You must be signed in to change notification settings - Fork 298
/
transformqueue.go
76 lines (61 loc) · 1.43 KB
/
transformqueue.go
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
package processor
import (
"container/heap"
"github.com/rudderlabs/rudder-server/processor/transformer"
)
type TransformRequestT struct {
Event []transformer.TransformerEventT
Stage string
ProcessingTime float64
Index int
}
type transformRequestPQ []*TransformRequestT
func (pq transformRequestPQ) Len() int {
return len(pq)
}
func (pq transformRequestPQ) Less(i, j int) bool {
return pq[i].ProcessingTime < pq[j].ProcessingTime
}
func (pq transformRequestPQ) Swap(i, j int) {
pq[i], pq[j] = pq[j], pq[i]
pq[i].Index = i
pq[j].Index = j
}
func (pq *transformRequestPQ) Push(x interface{}) {
n := len(*pq)
item := x.(*TransformRequestT)
item.Index = n
*pq = append(*pq, item)
}
func (pq *transformRequestPQ) Pop() interface{} {
old := *pq
n := len(old)
item := old[n-1]
item.Index = -1
*pq = old[0 : n-1]
return item
}
func (pq *transformRequestPQ) Top() *TransformRequestT {
item := (*pq)[0]
return item
}
func (pq *transformRequestPQ) Remove(item *TransformRequestT) {
heap.Remove(pq, item.Index)
}
func (pq *transformRequestPQ) Update(item, nextItem *TransformRequestT) {
index := item.Index
item = nextItem
item.Index = index
heap.Fix(pq, item.Index)
}
func (pq *transformRequestPQ) Add(item *TransformRequestT) {
heap.Push(pq, item)
}
func (pq *transformRequestPQ) RemoveTop() {
heap.Pop(pq)
}
func (pq *transformRequestPQ) Print() {
for _, v := range *pq {
pkgLogger.Debug(*v)
}
}