nodes

package
v0.8.0 Latest Latest
Warning

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

Go to latest
Published: Sep 23, 2026 License: Apache-2.0 Imports: 20 Imported by: 0

Documentation

Overview

Log nodes. Input is a map so a log node can be inserted anywhere in a pipeline without a projection: the previous node's output is logged as structured attributes and passed through unchanged.

Convention:

"message"  : string   — the log line (optional)
other keys : any      — structured attributes

Zero allocation when the level is disabled: slog.Default().Enabled is checked before any work is done.

Index

Constants

View Source
const (
	LogInfoName  = "log.info"
	LogWarnName  = "log.warn"
	LogErrorName = "log.error"
)

Variables

This section is empty.

Functions

func All added in v0.6.0

func All() []action.AnyAction

All returns every domain-neutral action this package provides.

Used by flow.Library() so a runtime that only needs the flow DSL can get a working action set without pulling in AI, sandbox, or provider dependencies.

func AssignPayload added in v0.8.0

func AssignPayload(input any, target any) error

AssignPayload decodes an input into target. It is the standard decode target used by every action that needs to feed a value into another action's request slot. Fast path for *any is deliberate: pipeline nodes almost always decode into `any`, and the type assertion avoids the reflect path in action.Assign.

func DistributeMaxItemsForTest added in v0.5.0

func DistributeMaxItemsForTest() int

func NewAssertAction added in v0.8.0

func NewAssertAction(condition, message string) (*action.BuiltAction[any, any], error)

NewAssertAction compiles an `assert(cond, msg)` guard. The condition is rewritten through PreprocessDotNotation so that leading-dot state references (`.game_over == false`, `.take >= 1`) compile identically to projections and loop conditions.

On a passing assertion the input payload is returned unchanged, so an assert can be inserted anywhere in a pipeline without altering the data flow.

func NewBenchCompareAction added in v0.5.0

func NewBenchCompareAction() action.AnyAction

func NewBenchRunAction added in v0.5.0

func NewBenchRunAction() action.AnyAction

NewBenchRunAction benchmarks another action. The target action is resolved against the registry that the flow compiler placed in the execution context, so this node works the same whether it is called directly or from inside a pipeline.

func NewBenchSaveAction added in v0.5.0

func NewBenchSaveAction() action.AnyAction

func NewDispatchAction added in v0.8.0

func NewDispatchAction() action.AnyAction

NewDispatchAction returns the dispatch node.

Dispatch order:

  1. Chosen, when it names a member.
  2. Fallback, when Chosen is empty or unknown.
  3. Every member in declaration order, first success wins.

The node fails only when every candidate fails and no fallback succeeded; the returned error wraps the last failure so the caller sees the actual cause, not a generic "dispatch failed".

Resolution of pool → members happens at call time, not at compile time, because the pool table is only known to the runner.

func NewDistributeMapAction added in v0.5.0

func NewDistributeMapAction() action.AnyAction

NewDistributeMapAction invokes another action once per item, bounded by concurrency. The target action is resolved against the registry placed in the execution context by the flow compiler.

func NewDistributeReduceAction added in v0.5.0

func NewDistributeReduceAction() action.AnyAction

func NewDynamicSaga

func NewDynamicSaga(name string, steps []SagaStep) *action.Builder[any, any]

NewDynamicSaga chains steps: each step's output feeds the next step's input. If a step fails, completed steps are compensated in LIFO order.

This is deliberately distinct from kernel/action.NewSaga, which runs every step against the same input. Here we need transformation chaining with rollback, which the flow DSL authors expect from `A -> B -> C` syntax.

func NewLogErrorAction added in v0.5.0

func NewLogErrorAction() action.AnyAction

func NewLogInfoAction added in v0.5.0

func NewLogInfoAction() action.AnyAction

func NewLogWarnAction added in v0.5.0

func NewLogWarnAction() action.AnyAction

func NewLoopAction

func NewLoopAction(
	bodyAction action.AnyAction,
	untilCondition string,
	maxTurns int,
) (*action.BuiltAction[any, any], error)

NewLoopAction compiles a `loop(body) until(cond)` construct. The until condition is rewritten through PreprocessDotNotation, so leading-dot state references work identically to projections and asserts.

func NewPanicAction added in v0.8.0

func NewPanicAction() action.AnyAction

NewPanicAction returns an action that panics. Use it in tests, guardrails, and escalation chains (`:onerror=["panic :msg=critical"]`).

A panic from this action is caught by BuiltAction.Do's recover() and converted to an error. In the DSL it is testable via `@onpanic`.

Distinct from `fail`:

  • fail → returns an error (recoverable, expected)
  • panic → unreachable state, invariant violation, bug

func NewProjectionAction

func NewProjectionAction(body string) (*action.BuiltAction[any, any], error)

NewProjectionAction compiles a projection body into a Kernel action.

The body is the inner content of a `{ ... }` projection in a flow pipeline, e.g. `attempt: attempt + 1, feedback: ...`. The parser (compiler/parseProjection) strips the outer braces before handing the body to this function.

Plain projection

A body without `...` produces a fresh object containing exactly the fields declared in the body:

{ to: .name, subject: "Welcome " + .tier }

Spread projection

A body that begins with `...` starts from the previous step's output and overrides only the fields explicitly listed:

{ ... }                          →  pass-through of the input
{ ..., attempt: attempt + 1 }    →  keep all fields, bump attempt

The spread token must be the first entry. Its output is exactly the input map with the listed fields overridden; no reserved names are added to the result.

Reserved identifiers

Inside the body the following identifiers are available in addition to the input's own fields:

__root__   — the raw, un-normalised input value
__state__  — the normalised object form of the input

`__state__` is what `...` expands to. Both are namespaced with `__` to avoid clashing with ordinary user data.

Array functions and the pipe operator

The projection body is evaluated by expr-lang and therefore has access to every built-in array function and to the pipe operator. The DSL layer does not restrict this surface; it only ensures that the expressions survive the two preprocessors (dot-notation rewrite and spread rewrite) unchanged.

Commonly used functions and shapes:

sortBy(array, #.field [, "asc"|"desc"])
groupBy(array, #.field)
filter(array, #.field == value)
map(array, #.field)
reduce(array, #acc + #.field, initial)
count(array [, #.field == value])
uniq(array)
flatten(array)
concat(a, b, ...)
first(array)
last(array)
take(array, n)
reverse(array)
toJSON(value)
fromJSON(string)
len(collection)

The predicate scope marker is `#`, which refers to the current element. The dot-notation preprocessor preserves `#.field` exactly as written; a bare `.field` outside a predicate is rewritten to `field` to match the flattened environment produced by NormalizeEnv.

A realistic chained projection:

{
  ...,
  critical: filter(findings, #.severity == "critical"),
  top:      findings | sortBy(#.line, "desc") | take(10),
  count:    count(findings, #.severity == "warning"),
}

Nil safety

Two operators from expr-lang are especially useful for optional fields:

feedback ?? "none"       →  fallback when feedback is nil
patch?.source_code       →  safe field access on a possibly nil object

func NewPromptNode

func NewPromptNode(cfg PromptConfig) action.AnyAction

NewPromptNode creates a reusable, pre-configured prompt template action.

func NewSupervisorNode

func NewSupervisorNode(name string) action.AnyAction

NewSupervisorNode builds a supervisor that compiles and runs child pipelines on the fly. The compiler is obtained from the execution context, so the node works inside any flow without constructor args.

func NormalizeEnv added in v0.8.0

func NormalizeEnv(input any) any

NormalizeEnv converts a nested value into a flat lookup map that can be handed to expr-lang as its environment.

Behavior

  • Maps are copied. Every key is exposed three ways: verbatim, in lower-case, and — when the key itself contains dots — also as a nested path. `{"user.name": "Alice"}` therefore yields both `user.name` (literal key) and `user.name` reachable as `user["name"]` or `user.name` after PreprocessDotNotation.

  • Structs are flattened by their exported field names, their JSON tags, and their lower-cased Go names. This lets expressions use any of the three spellings interchangeably.

  • Pointers and interfaces are dereferenced. Nil values become nil.

  • All other values are returned unchanged.

Cost

One map allocation per nested map or struct. Small, predictable, and bounded by the shape of the input.

func PreprocessDotNotation added in v0.8.0

func PreprocessDotNotation(src string) (out string, usesRoot bool)

PreprocessDotNotation rewrites a projection expression so that the leading-dot shorthand can be used with flattened state maps.

Rewrites

.foo              →  foo                     (identifier shorthand)
.foo.bar          →  foo.bar
.foo[0]           →  foo[0]
.foo.bar[2].baz   →  foo.bar[2].baz

The leading-dot shorthand is *not* rewritten when:

  • the dot is part of a floating-point literal (".5", "3.14"),
  • the dot follows an identifier, ")", or "]", because in that case it is an ordinary member access on an existing value, and
  • the dot follows a predicate scope marker `#`, which is the modern expr-lang syntax for referring to the current element of a built-in array function such as `sortBy`, `groupBy`, `map`, `filter`, or `reduce`.

In the last case the dot must be preserved, otherwise the predicate body would be corrupted (`#.Age` must not become `#Age`).

Root reference

When the expression contains a bare `.` that is not followed by an identifier, the dot is replaced with the reserved identifier `__root__` and `usesRoot` is set to true. This allows a projection to reference the whole previous state as `__root__` explicitly.

The function is intentionally lexical: it does not parse the expression, it only classifies dots by their immediate neighbours. This is sufficient for the projections produced by the DSL and keeps the hot path allocation-free apart from the builder itself.

func PreprocessSpread added in v0.8.0

func PreprocessSpread(body string) (string, bool, error)

PreprocessSpread prepares a projection body for expr-lang.

Three shapes are supported:

  1. Bare expression. A body that does not look like a map literal body is returned unchanged. This is what allows function calls and pipelines to be used in a projection slot:

    sortBy(findings, #.severity) findings | filter(#.ok) | take(10) count(items, #.active)

  2. Map body without spread. A body that starts with `key:` (or `"key":`, `'key':`, “ `key`: “) is wrapped in a map literal:

    a: 1, b: 2 → { a: 1, b: 2 }

  3. Spread. A body that begins with `...` starts from the previous step's state and overrides only the listed fields:

    ... → __state__ ..., attempt: attempt + 1 → nexss_spread_merge(__state__, { attempt: attempt + 1 })

The spread token must be the first entry. A misplaced `...` produces a compile error with a helpful message rather than a cryptic expr-lang parse error.

Types

type BenchCompareReq added in v0.5.0

type BenchCompareReq struct {
	Baseline     string  `json:"baseline" validate:"required" usage:"Baseline JSON path (validated by xfs.Rel)"`
	TolerancePct float64 `json:"tolerance_pct,omitempty"      usage:"Allowed regression percent (default 5)"`
	BenchRunRes
}

type BenchCompareRes added in v0.5.0

type BenchCompareRes struct {
	Baseline    string                 `json:"baseline"`
	Pass        bool                   `json:"pass"`
	Regressions []string               `json:"regressions,omitempty"`
	Metrics     map[string]BenchMetric `json:"metrics"`
}

type BenchMetric added in v0.5.0

type BenchMetric struct {
	Baseline  float64 `json:"baseline"`
	Current   float64 `json:"current"`
	DeltaPct  float64 `json:"delta_pct"`
	Regressed bool    `json:"regressed"`
}

type BenchRunReq added in v0.5.0

type BenchRunReq struct {
	Action     string         `json:"action"               validate:"required" usage:"Node name to benchmark"`
	Iterations int            `json:"iterations,omitempty"                    usage:"Timed iterations (default 50)"`
	Warmup     int            `json:"warmup,omitempty"                        usage:"Discarded pre-runs (default 3)"`
	Payload    map[string]any `json:"payload,omitempty"                       usage:"Request passed to each invocation"`
}

type BenchRunRes added in v0.5.0

type BenchRunRes struct {
	Action     string  `json:"action"`
	Iterations int     `json:"iterations"`
	Warmup     int     `json:"warmup"`
	Errors     int     `json:"errors"`
	MinMs      float64 `json:"min_ms"`
	MaxMs      float64 `json:"max_ms"`
	MeanMs     float64 `json:"mean_ms"`
	P50Ms      float64 `json:"p50_ms"`
	P95Ms      float64 `json:"p95_ms"`
	P99Ms      float64 `json:"p99_ms"`
	RPS        float64 `json:"rps"`
	ElapsedMs  int64   `json:"elapsed_ms"`
}

type BenchSaveReq added in v0.5.0

type BenchSaveReq struct {
	File string `json:"file" validate:"required" usage:"Relative path to write (validated by xfs.Rel)"`
	BenchRunRes
}

type BenchSaveRes added in v0.5.0

type BenchSaveRes struct {
	File  string `json:"file"`
	Bytes int    `json:"bytes"`
	BenchRunRes
}

type ChildResult

type ChildResult struct {
	TaskID   string `json:"task_id"`
	Output   any    `json:"output,omitempty"`
	Error    string `json:"error,omitempty"`
	Duration int64  `json:"duration_ms"`
}

type ChildTask

type ChildTask struct {
	ID        string `json:"id"`
	DSL       string `json:"dsl"`
	Payload   any    `json:"payload"`
	TimeoutMS int64  `json:"timeout_ms,omitempty"`
}

type DispatchReq added in v0.8.0

type DispatchReq struct {
	Pool     string `json:"pool,omitempty"`
	Members  string `json:"members,omitempty"`
	Chosen   string `json:"chosen,omitempty"`
	Fallback string `json:"fallback,omitempty"`
	Payload  any    `json:"payload,omitempty"`
}

DispatchReq is the request shape for the dispatch node.

Members is a comma-separated list of action names, or empty when Pool is set. Pool names a @pool declaration, which the node expands into Members by reading the pool table from the execution context.

Chosen is the name the upstream node or LLM selected. Fallback is the name used when Chosen is missing or fails. Payload is what every candidate is invoked with.

type DistributeItem added in v0.5.0

type DistributeItem struct {
	OK     bool   `json:"ok"`
	Result any    `json:"result,omitempty"`
	Error  string `json:"error,omitempty"`
}

type DistributeMapReq added in v0.5.0

type DistributeMapReq struct {
	Action      string `json:"action"      validate:"required" usage:"Node name to invoke for each item"`
	Concurrency int    `json:"concurrency,omitempty"           usage:"Max in-flight invocations (default 4, max 256)"`
	Items       []any  `json:"items"       validate:"required" usage:"Items to fan out. Passed unchanged to Action."`
}

type DistributeMapRes added in v0.5.0

type DistributeMapRes struct {
	Items     []DistributeItem `json:"items"`
	Succeeded int              `json:"succeeded"`
	Failed    int              `json:"failed"`
}

type DistributeReduceReq added in v0.5.0

type DistributeReduceReq struct {
	Strategy string           `json:"strategy" validate:"required,oneof=collect all_pass any_pass first_success count"`
	Items    []DistributeItem `json:"items"    validate:"required"`
}

type DistributeReduceRes added in v0.5.0

type DistributeReduceRes struct {
	Strategy  string `json:"strategy"`
	AllPass   bool   `json:"all_pass,omitempty"`
	AnyPass   bool   `json:"any_pass,omitempty"`
	Succeeded int    `json:"succeeded,omitempty"`
	Failed    int    `json:"failed,omitempty"`
	Result    any    `json:"result,omitempty"`
	Results   []any  `json:"results,omitempty"`
}

type PanicReq added in v0.8.0

type PanicReq struct {
	Msg string `json:"msg" cli:"msg" usage:"Message passed to panic()"`
}

PanicReq is the request DTO for the `panic` action.

type PromptConfig

type PromptConfig struct {
	Name          string         `json:"name"`
	Description   string         `json:"description"`
	SystemPrompt  string         `json:"system_prompt"`
	UserTemplate  string         `json:"user_template"`
	DefaultParams map[string]any `json:"default_params,omitempty"`
	Timeout       time.Duration  `json:"timeout,omitempty"`
}

type PromptReq

type PromptReq struct {
	Input  any            `json:"input,omitempty"`
	Params map[string]any `json:"params,omitempty"`
	Prompt string         `json:"prompt,omitempty"`
}

type PromptRes

type PromptRes struct {
	SystemPrompt string         `json:"system_prompt"`
	RenderedUser string         `json:"rendered_user"`
	Params       map[string]any `json:"params"`
}

type SagaStep

type SagaStep struct {
	NodeID     string
	Forward    action.AnyAction
	Compensate action.AnyAction
}

type SupervisorReq

type SupervisorReq struct {
	Tasks []ChildTask `json:"tasks" validate:"required"`
}

type SupervisorRes

type SupervisorRes struct {
	Total     int           `json:"total"`
	Succeeded int           `json:"succeeded"`
	Failed    int           `json:"failed"`
	Results   []ChildResult `json:"results"`
}

Directories

Path Synopsis
flow/nodes/fsio/library.go
flow/nodes/fsio/library.go

Jump to

Keyboard shortcuts

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