forked from elastic/beats
-
Notifications
You must be signed in to change notification settings - Fork 0
/
out.go
91 lines (71 loc) · 1.96 KB
/
out.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
package stress
import (
"math/rand"
"time"
"github.com/elastic/beats/libbeat/beat"
"github.com/elastic/beats/libbeat/common"
"github.com/elastic/beats/libbeat/outputs"
"github.com/elastic/beats/libbeat/publisher"
)
type testOutput struct {
config testOutputConfig
observer outputs.Observer
batchCount int
}
type testOutputConfig struct {
Worker int `config:"worker" validate:"min=1"`
BulkMaxSize int `config:"bulk_max_size"`
Retry int `config:"retry"`
MinWait time.Duration `config:"min_wait"`
MaxWait time.Duration `config:"max_wait"`
Fail struct {
EveryBatch int
}
}
var defaultTestOutputConfig = testOutputConfig{
Worker: 1,
BulkMaxSize: 64,
}
func init() {
outputs.RegisterType("test", makeTestOutput)
}
func makeTestOutput(beat beat.Info, observer outputs.Observer, cfg *common.Config) (outputs.Group, error) {
config := defaultTestOutputConfig
if err := cfg.Unpack(&config); err != nil {
return outputs.Fail(err)
}
clients := make([]outputs.Client, config.Worker)
for i := range clients {
client := &testOutput{config: config, observer: observer}
clients[i] = client
}
return outputs.Success(config.BulkMaxSize, config.Retry, clients...)
}
func (*testOutput) Close() error { return nil }
func (t *testOutput) Publish(batch publisher.Batch) error {
config := &t.config
n := len(batch.Events())
t.observer.NewBatch(n)
min := int64(config.MinWait)
max := int64(config.MaxWait)
if max > 0 && min < max {
waitFor := rand.Int63n(max-min) + min
// TODO: make wait interruptable via `Close`
time.Sleep(time.Duration(waitFor))
}
// fail complete batch
if config.Fail.EveryBatch > 0 {
t.batchCount++
if config.Fail.EveryBatch == t.batchCount {
t.batchCount = 0
t.observer.Failed(n)
batch.Retry()
return nil
}
}
// TODO: add support to fail single events at end of batch or randomly
// ack complete batch
batch.ACK()
t.observer.Acked(n)
return nil
}