-
Notifications
You must be signed in to change notification settings - Fork 2
/
Copy pathpipeline.go
140 lines (125 loc) · 3.63 KB
/
pipeline.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
140
// Copyright (C) ENEO Tecnologia SL - 2016
//
// Authors: Diego Fernández Barrera <dfernandez@redborder.com> <bigomby@gmail.com>
// Eugenio Pérez Martín <eugenio@redborder.com>
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Lesser General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// This program is distributed in the hope that it will be useful,
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Lesser General Public License for more details.
//
// You should have received a copy of the GNU General Public License
// along with this program. If not, see <https://www.gnu.org/licenses/lgpl-3.0.txt>.
package rbforwarder
import (
"github.com/oleiade/lane"
"github.com/redBorder/rbforwarder/utils"
)
// Component contains information about a pipeline component
type Component struct {
composser utils.Composer
pool chan chan *utils.Message
}
// pipeline contains the components
type pipeline struct {
components []Component
input chan *utils.Message
retry chan *utils.Message
output chan *utils.Message
}
// newPipeline creates a new Backend
func newPipeline(input, retry, output chan *utils.Message) *pipeline {
return &pipeline{
input: input,
retry: retry,
output: output,
}
}
// PushComponent adds a new component to the pipeline
func (p *pipeline) PushComponent(composser utils.Composer) {
p.components = append(p.components, struct {
composser utils.Composer
pool chan chan *utils.Message
}{
composser: composser,
pool: make(chan chan *utils.Message, composser.Workers()),
})
}
func (p *pipeline) Run() {
for index, component := range p.components {
for w := 0; w < component.composser.Workers(); w++ {
go func(w, index int, component Component) {
component.composser = component.composser.Spawn(w)
messages := make(chan *utils.Message)
component.pool <- messages
for m := range messages {
component.composser.OnMessage(m,
// Done function
func(m *utils.Message, code int, status string) {
// If there is another component next in the pipeline send the
// messate to it. I other case send the message to the report
// handler
if code == 0 && len(p.components)-1 > index {
nextWorker := <-p.components[index+1].pool
nextWorker <- m
} else {
reports := lane.NewStack()
for !m.Reports.Empty() {
rep := m.Reports.Pop().(Report)
rep.Component = index
rep.Code = code
rep.Status = status
reports.Push(rep)
}
m.Reports = reports
p.output <- m
}
})
component.pool <- messages
}
}(w, index, component)
}
}
go func() {
for {
select {
case m, ok := <-p.retry:
if ok {
rep := m.Reports.Head().(Report)
worker := <-p.components[rep.Component].pool
worker <- m
}
continue
default:
}
select {
case m, ok := <-p.retry:
if ok {
rep := m.Reports.Head().(Report)
worker := <-p.components[rep.Component].pool
worker <- m
}
continue
case m, ok := <-p.input:
// If input channel has been closed, close output channel
if !ok {
for _, component := range p.components {
for i := 0; i < component.composser.Workers(); i++ {
worker := <-component.pool
close(worker)
}
}
close(p.output)
} else {
worker := <-p.components[0].pool
worker <- m
}
}
}
}()
}