-
Notifications
You must be signed in to change notification settings - Fork 152
/
source.go
38 lines (30 loc) · 920 Bytes
/
source.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
package execute
import (
"context"
"fmt"
"github.com/influxdata/flux"
"github.com/influxdata/flux/plan"
)
type Node interface {
AddTransformation(t Transformation)
}
// MetadataNode is a node that has additional metadata
// that should be added to the result after it is
// processed.
type MetadataNode interface {
Node
Metadata() flux.Metadata
}
type Source interface {
Node
Run(ctx context.Context)
}
type CreateSource func(spec plan.ProcedureSpec, id DatasetID, ctx Administration) (Source, error)
type CreateNewPlannerSource func(spec plan.ProcedureSpec, id DatasetID, ctx Administration) (Source, error)
var procedureToSource = make(map[plan.ProcedureKind]CreateNewPlannerSource)
func RegisterSource(k plan.ProcedureKind, c CreateNewPlannerSource) {
if procedureToSource[k] != nil {
panic(fmt.Errorf("duplicate registration for source with procedure kind %v", k))
}
procedureToSource[k] = c
}