nodes_supervisor

package
v0.18.0 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 7, 2026 License: Apache-2.0 Imports: 8 Imported by: 0

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

func Bundle

func Bundle(_ map[string]string) core.Bundle

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.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL