(seq dag.Seq, post func(dag.Seq) (dag.Seq, error))
| 90 | } |
| 91 | |
| 92 | func walkEntries(seq dag.Seq, post func(dag.Seq) (dag.Seq, error)) (dag.Seq, error) { |
| 93 | for _, op := range seq { |
| 94 | switch op := op.(type) { |
| 95 | case *dag.ForkOp: |
| 96 | for k := range op.Paths { |
| 97 | seq, err := walkEntries(op.Paths[k], post) |
| 98 | if err != nil { |
| 99 | return nil, err |
| 100 | } |
| 101 | op.Paths[k] = seq |
| 102 | } |
| 103 | case *dag.ScatterOp: |
| 104 | for k := range op.Paths { |
| 105 | seq, err := walkEntries(op.Paths[k], post) |
| 106 | if err != nil { |
| 107 | return nil, err |
| 108 | } |
| 109 | op.Paths[k] = seq |
| 110 | } |
| 111 | } |
| 112 | } |
| 113 | return post(seq) |
| 114 | } |
| 115 | |
| 116 | // Optimize transforms the DAG by attempting to lift stateless operators |
| 117 | // from the downstream sequence into the trunk of each data source in the From |
no outgoing calls
no test coverage detected