Documentation
¶
Overview ¶
Package nodes_supervisor provides the supervisor action: it compiles and runs a set of named child pipelines concurrently, isolating each child's panic, error, and timeout.
Typical use:
{ tasks: [
{ id: "a", dsl: "noop", payload: { x: 1 } },
{ id: "b", dsl: "log.info", timeout_ms: 500 }
]} -> supervisor.run
Index ¶
Constants ¶
View Source
const ID = "nodes_supervisor"
Variables ¶
View Source
var Supervisor = action.New("supervisor.run", func(ctx context.Context, req SupervisorReq) (SupervisorRes, error) { if len(req.Tasks) == 0 { return SupervisorRes{}, xerr.BadRequest("supervisor: no child tasks provided") } compiler := contracts.CompilerFromContext(ctx) if compiler == nil { return SupervisorRes{}, xerr.Internal("supervisor: no pipeline compiler in context") } results := make([]ChildResult, len(req.Tasks)) group, groupCtx := errgroup.WithContext(ctx) for i, task := range req.Tasks { index, current := i, task group.Go(func() error { results[index] = runChild(groupCtx, compiler, current) return nil }) } if err := group.Wait(); err != nil { return SupervisorRes{}, err } var succeeded, failed int for i := range results { r := &results[i] if r.Error != "" { failed++ } else { succeeded++ } } return SupervisorRes{ Total: len(results), Succeeded: succeeded, Failed: failed, Results: results, }, nil }).Description("Compile and run named child pipelines concurrently"). Build()
Supervisor compiles and runs each task's DSL. All children run concurrently; each child is isolated (panic, error, timeout). The compiler is read from the execution context.
Functions ¶
Types ¶
type ChildResult ¶
type ChildResult struct {
TaskID string `json:"task_id"`
Output any `json:"output,omitempty"`
Error string `json:"error,omitempty"`
DurationMs int64 `json:"duration_ms"`
}
ChildResult reports one child's outcome. Error is a plain string so the JSON shape stays stable across error kinds.
type ChildTask ¶
type ChildTask struct {
ID string `json:"id"`
DSL string `json:"dsl"`
Payload any `json:"payload"`
TimeoutMS int64 `json:"timeout_ms,omitempty"`
}
ChildTask describes one pipeline to compile and run.
type SupervisorReq ¶
type SupervisorReq struct {
Tasks []ChildTask `json:"tasks" validate:"required"`
}
SupervisorReq is the set of children to run.
type SupervisorRes ¶
type SupervisorRes struct {
Total int `json:"total"`
Succeeded int `json:"succeeded"`
Failed int `json:"failed"`
Results []ChildResult `json:"results"`
}
SupervisorRes aggregates the outcome of every child.
Click to show internal directories.
Click to hide internal directories.