/
builtin_fn_flow.go
125 lines (110 loc) · 2.35 KB
/
builtin_fn_flow.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
package eval
import (
"errors"
"sync"
"github.com/elves/elvish/util"
)
// Flow control.
func init() {
addBuiltinFns(map[string]interface{}{
"run-parallel": runParallel,
// Exception and control
"fail": fail,
"multi-error": multiErrorFn,
"return": returnFn,
"break": breakFn,
"continue": continueFn,
// Iterations.
"each": each,
"peach": peach,
})
}
func runParallel(fm *Frame, functions ...Callable) error {
var waitg sync.WaitGroup
waitg.Add(len(functions))
exceptions := make([]*Exception, len(functions))
for i, function := range functions {
go func(fm2 *Frame, function Callable, exception **Exception) {
err := fm2.Call(function, NoArgs, NoOpts)
if err != nil {
*exception = err.(*Exception)
}
waitg.Done()
}(fm.fork("[run-parallel function]"), function, &exceptions[i])
}
waitg.Wait()
return ComposeExceptionsFromPipeline(exceptions)
}
// each takes a single closure and applies it to all input values.
func each(fm *Frame, f Callable, inputs Inputs) error {
broken := false
var err error
inputs(func(v interface{}) {
if broken {
return
}
newFm := fm.fork("closure of each")
newFm.ports[0] = DevNullClosedChan
ex := newFm.Call(f, []interface{}{v}, NoOpts)
newFm.Close()
if ex != nil {
switch ex.(*Exception).Cause {
case nil, Continue:
// nop
case Break:
broken = true
default:
broken = true
err = ex
}
}
})
return err
}
// peach takes a single closure and applies it to all input values in parallel.
func peach(fm *Frame, f Callable, inputs Inputs) error {
var w sync.WaitGroup
broken := false
var err error
inputs(func(v interface{}) {
if broken || err != nil {
return
}
w.Add(1)
go func() {
newFm := fm.fork("closure of peach")
newFm.ports[0] = DevNullClosedChan
ex := newFm.Call(f, []interface{}{v}, NoOpts)
newFm.Close()
if ex != nil {
switch ex.(*Exception).Cause {
case nil, Continue:
// nop
case Break:
broken = true
default:
broken = true
err = util.Errors(err, ex)
}
}
w.Done()
}()
})
w.Wait()
return err
}
func fail(msg string) error {
return errors.New(msg)
}
func multiErrorFn(excs ...*Exception) error {
return PipelineError{excs}
}
func returnFn() error {
return Return
}
func breakFn() error {
return Break
}
func continueFn() error {
return Continue
}