forked from redpanda-data/connect
-
Notifications
You must be signed in to change notification settings - Fork 0
/
type.go
108 lines (91 loc) · 2.34 KB
/
type.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
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
package checkpoint
import (
"sync"
)
// Type keeps track of a sequence of pending checkpoint payloads, and as pending
// checkpoints are resolved it retains the latest fully resolved payload in the
// sequence where all prior sequence checkpoints are also resolved.
//
// Also keeps track of the logical size of the unresolved sequence, which allows
// for limiting the number of pending checkpoints.
type Type struct {
positionOffset int64
checkpoint interface{}
latest, earliest *node
}
// New returns a new checkpointer.
func New() *Type {
return &Type{}
}
// Track a new unresolved payload. This payload will be cached until it is
// marked as resolved. While it is cached no more recent payload will ever be
// committed.
//
// While the returned resolve funcs can be called from any goroutine, it
// is assumed that Track is called from a single goroutine.
func (t *Type) Track(payload interface{}, batchSize int64) func() interface{} {
newNode := getNode()
newNode.payload = payload
newNode.position = batchSize
if t.earliest == nil {
t.earliest = newNode
}
if t.latest != nil {
newNode.prev = t.latest
newNode.position += t.latest.position
t.latest.next = newNode
}
t.latest = newNode
return func() interface{} {
if newNode.prev != nil {
newNode.prev.position = newNode.position
newNode.prev.payload = newNode.payload
newNode.prev.next = newNode.next
} else {
t.checkpoint = newNode.payload
t.positionOffset = newNode.position
t.earliest = newNode.next
}
if newNode.next != nil {
newNode.next.prev = newNode.prev
} else {
t.latest = newNode.prev
if t.latest == nil {
t.positionOffset = 0
}
}
putNode(newNode)
return t.checkpoint
}
}
// Pending returns the gap between the earliest and latests unresolved messages.
func (t *Type) Pending() int64 {
if t.latest == nil {
return 0
}
return t.latest.position - t.positionOffset
}
// Highest returns the payload of the highest resolved checkpoint.
func (t *Type) Highest() interface{} {
return t.checkpoint
}
type node struct {
position int64
payload interface{}
prev, next *node
}
var nodePool = &sync.Pool{
New: func() interface{} {
return &node{}
},
}
func getNode() *node {
return nodePool.Get().(*node)
}
func putNode(node *node) {
node.position = 0
node.payload = nil
node.prev = nil
node.next = nil
nodePool.Put(node)
}