forked from elastic/beats
-
Notifications
You must be signed in to change notification settings - Fork 0
/
batch.go
71 lines (57 loc) · 1.34 KB
/
batch.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
package outest
import (
"github.com/elastic/beats/libbeat/beat"
"github.com/elastic/beats/libbeat/publisher"
)
type Batch struct {
events []publisher.Event
Signals []BatchSignal
OnSignal func(sig BatchSignal)
}
type BatchSignal struct {
Tag BatchSignalTag
Events []publisher.Event
}
type BatchSignalTag uint8
const (
BatchACK BatchSignalTag = iota
BatchDrop
BatchRetry
BatchRetryEvents
BatchCancelled
BatchCancelledEvents
)
func NewBatch(in ...beat.Event) *Batch {
events := make([]publisher.Event, len(in))
for i, c := range in {
events[i] = publisher.Event{Content: c}
}
return &Batch{events: events}
}
func (b *Batch) Events() []publisher.Event {
return b.events
}
func (b *Batch) ACK() {
b.doSignal(BatchSignal{Tag: BatchACK})
}
func (b *Batch) Drop() {
b.doSignal(BatchSignal{Tag: BatchDrop})
}
func (b *Batch) Retry() {
b.doSignal(BatchSignal{Tag: BatchRetry})
}
func (b *Batch) RetryEvents(events []publisher.Event) {
b.doSignal(BatchSignal{Tag: BatchRetryEvents, Events: events})
}
func (b *Batch) Cancelled() {
b.doSignal(BatchSignal{Tag: BatchCancelled})
}
func (b *Batch) CancelledEvents(events []publisher.Event) {
b.doSignal(BatchSignal{Tag: BatchCancelledEvents, Events: events})
}
func (b *Batch) doSignal(sig BatchSignal) {
b.Signals = append(b.Signals, sig)
if b.OnSignal != nil {
b.OnSignal(sig)
}
}