forked from Roman2K/scat
-
Notifications
You must be signed in to change notification settings - Fork 0
/
chain.go
112 lines (99 loc) · 1.96 KB
/
chain.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
package procs
import (
"sync"
"github.com/Roman2K/scat"
)
type Chain []Proc
var _ Proc = Chain{}
func (chain Chain) Process(c *scat.Chunk) <-chan Res {
procs := chain
enders := chain.endProcs()
if len(enders) > 0 {
ecp := endCallProc{chunk: c, enders: enders}
newProcs := make([]Proc, len(procs)+1)
copy(newProcs, procs)
newProcs[len(newProcs)-1] = ecp
procs = newProcs
}
in := SingleRes(c, nil)
var out chan Res
for _, proc := range procs {
out = make(chan Res)
go process(out, in, proc)
in = out
}
if out == nil {
return in
}
return out
}
func (procs Chain) endProcs() (enders []EndProc) {
for _, p := range procs {
if e, ok := underlying(p).(EndProc); ok {
enders = append(enders, e)
}
}
return
}
func (procs Chain) Finish() error {
return finishFuncs(procs).FirstErr()
}
func process(out chan<- Res, in <-chan Res, proc Proc) {
defer close(out)
wg := sync.WaitGroup{}
for res := range in {
var ch <-chan Res
if res.Err != nil {
if errp, ok := underlying(proc).(ErrProc); ok && res.Chunk != nil {
ch = errp.ProcessErr(res.Chunk, res.Err)
} else {
out <- res
continue
}
} else {
ch = proc.Process(res.Chunk)
}
wg.Add(1)
go func() {
defer wg.Done()
for res := range ch {
out <- res
}
}()
}
wg.Wait()
if ecp, ok := proc.(endCallProc); ok {
err := ecp.processEnd()
if err != nil {
out <- Res{Err: err}
}
}
}
type endCallProc struct {
chunk *scat.Chunk
enders []EndProc
}
func (ecp endCallProc) Process(c *scat.Chunk) <-chan Res {
return InplaceFunc(ecp.process).Process(c)
}
func (ecp endCallProc) process(final *scat.Chunk) (err error) {
for _, ender := range ecp.enders {
err = ender.ProcessFinal(ecp.chunk, final)
if err != nil {
return
}
}
return
}
func (ecp endCallProc) processEnd() (err error) {
for _, ender := range ecp.enders {
err = ender.ProcessEnd(ecp.chunk)
if err != nil {
return
}
}
return
}
func (ecp endCallProc) Finish() error {
return nil
}