-
Notifications
You must be signed in to change notification settings - Fork 3
/
rate_synchronizer.go
139 lines (122 loc) · 3.31 KB
/
rate_synchronizer.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
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
package steps
import (
"errors"
"sync"
"github.com/antongulenko/golib"
"github.com/bitflow-stream/go-bitflow/bitflow"
"github.com/bitflow-stream/go-bitflow/script/reg"
)
type PipelineRateSynchronizer struct {
steps []*synchronizationStep
startOnce sync.Once
ChannelSize int
ChannelCloseHook func(lastSample *bitflow.Sample, lastHeader *bitflow.Header)
}
func RegisterPipelineRateSynchronizer(b reg.ProcessorRegistry) {
synchronization_keys := make(map[string]*PipelineRateSynchronizer)
create := func(p *bitflow.SamplePipeline, params map[string]interface{}) error {
chanSize := params["buf"].(int)
key := params["key"].(string)
synchronizer, ok := synchronization_keys[key]
if !ok {
synchronizer = &PipelineRateSynchronizer{
ChannelSize: chanSize,
}
synchronization_keys[key] = synchronizer
} else if synchronizer.ChannelSize != chanSize {
return reg.ParameterError("buf", errors.New("synchronize() steps with the same 'key' parameter must all have the same 'buf' parameter"))
}
p.Add(synchronizer.NewSynchronizationStep())
return nil
}
b.RegisterStep("synchronize", create,
"Synchronize the number of samples going through each synchronize() step with the same key parameter").
Required("key", reg.String()).
Optional("buf", reg.Int(), 5)
}
func (s *PipelineRateSynchronizer) NewSynchronizationStep() bitflow.SampleProcessor {
chanSize := s.ChannelSize
if chanSize < 1 {
chanSize = 1
}
step := &synchronizationStep{
synchronizer: s,
queue: make(chan bitflow.SampleAndHeader, chanSize),
running: true,
}
s.steps = append(s.steps, step)
return step
}
func (s *PipelineRateSynchronizer) start(wg *sync.WaitGroup) {
s.startOnce.Do(func() {
wg.Add(1)
go s.process(wg)
})
}
func (s *PipelineRateSynchronizer) process(wg *sync.WaitGroup) {
defer wg.Done()
for {
runningSteps := 0
for _, step := range s.steps {
if step.running {
if sample := <-step.queue; sample.Sample != nil {
step.outputSample(sample)
runningSteps++
} else {
step.running = false
step.CloseSink()
}
}
}
if runningSteps == 0 {
break
}
}
}
type synchronizationStep struct {
bitflow.NoopProcessor
synchronizer *PipelineRateSynchronizer
queue chan bitflow.SampleAndHeader
running bool
closeSinkOnce sync.Once
err error
lastSample *bitflow.Sample
lastHeader *bitflow.Header
}
func (s *synchronizationStep) String() string {
return "Synchronize processing rate"
}
func (s *synchronizationStep) Start(wg *sync.WaitGroup) golib.StopChan {
s.synchronizer.start(wg)
return s.NoopProcessor.Start(wg)
}
func (s *synchronizationStep) Sample(sample *bitflow.Sample, header *bitflow.Header) error {
if e := s.err; e != nil {
s.err = nil
return e
}
s.lastSample = sample
s.lastHeader = header
s.queue <- bitflow.SampleAndHeader{
Sample: sample,
Header: header,
}
return nil
}
func (s *synchronizationStep) Close() {
close(s.queue)
}
func (s *synchronizationStep) CloseSink() {
s.closeSinkOnce.Do(func() {
if c := s.synchronizer.ChannelCloseHook; c != nil {
c(s.lastSample, s.lastHeader)
}
s.NoopProcessor.CloseSink()
})
}
func (s *synchronizationStep) outputSample(sample bitflow.SampleAndHeader) {
err := s.NoopProcessor.Sample(sample.Sample, sample.Header)
if err != nil && s.err == nil {
s.err = err
}
}