/
builder.go
82 lines (76 loc) · 1.84 KB
/
builder.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
package dag
import (
"context"
"golang.org/x/sync/errgroup"
"k8s.io/klog/v2"
)
func BuildDag(nodes []*DagNode) *Dag {
dag := &Dag{
indegree: make(map[string]int, len(nodes)),
length: len(nodes),
raw: make(map[string]*DagNode, len(nodes)),
}
for i, node := range nodes {
dag.raw[node.Name] = nodes[i]
dag.indegree[node.Name] = len(node.Depends)
}
for name, indegree := range dag.indegree {
if indegree == 0 {
dag.raw[name].PreContext = make(DagContext)
dag.head = append(dag.head, dag.raw[name])
}
for _, dependName := range dag.raw[name].Depends {
dag.raw[dependName].Next = append(dag.raw[dependName].Next,
dag.raw[name],
)
}
}
return dag
}
func (dag *Dag) Run(ctx context.Context) error {
dag.finished = 0
err := dag.executeIndgreeZeroNodes(ctx, dag.head)
klog.Infof("dag task finished %d tasks, and left %d tasks not to execute", dag.finished, dag.length-dag.finished)
return err
}
func (dag *Dag) executeIndgreeZeroNodes(ctx context.Context, nodes []*DagNode) error {
if len(nodes) == 0 {
return nil
}
if err := dag.executeNodes(ctx, nodes); err != nil {
return err
}
dag.finished += len(nodes)
tmpNodes := []*DagNode{}
for _, head := range nodes {
for _, next := range head.Next {
dag.indegree[next.Name]--
if dag.indegree[next.Name] == 0 {
tmpNodes = append(tmpNodes, next)
}
}
}
return dag.executeIndgreeZeroNodes(ctx, tmpNodes)
}
func (dag *Dag) executeNodes(ctx context.Context, nodes []*DagNode) error {
wg := errgroup.Group{}
exec := func(nd *DagNode) func() error {
return func() error {
defer func() {
nd.done = true
}()
// wait depends
nd.WaitDependDone(ctx.Done(), dag.raw)
// execute
if !nd.ConditionValid() {
// not execute
return nil
}
return nd.Execute()
}
}
for _, node := range nodes {
wg.Go(exec(node))
}
return wg.Wait()
}