Documentation
¶
Overview ¶
Package flow — compiler.go owns the Compiler type and its top-level Compile method. Supporting code lives in sibling files:
compiler_ctx.go context keys and helpers used by Compile compiler_hooks.go profile hook resolution compiler_resolve.go capability lookup and payload unpacking
The split keeps this file to the compilation pipeline itself, so the running commentary on Compile stays readable even as the rest of the package grows.
flow/cost.go
Index ¶
- Constants
- Variables
- func AcquireStateFromGraphState(s *State) *dag.State
- func AsCostHook(reserver cost.Reserver, estimateMicros int64, _ ...int64) action.AnyHook
- func BaseLibrary() action.Library
- func BuildCatalogAction(registry *action.Registry, bindings ...action.Binding) action.AnyAction
- func BuildResolvedConfig(pre *core.Preprocessed, cliArgs []string) map[string]string
- func CoerceLiteralValue(trimmed string) any
- func CompilePipeline(expr string, reg *action.Registry, opts ...CompileOption) (*action.Builder[any, any], error)
- func CompileSaga(expr string, reg *action.Registry) (*action.Builder[any, any], error)
- func ContributeRegistry(pre *core.Preprocessed, base *action.Registry) (*action.Registry, error)
- func CountStreamResult(res any) int
- func EvaluateCondition(condition string, state *State) (bool, error)
- func Execute[Req, Res any](ctx context.Context, act *action.BuiltAction[Req, Res], req Req) (Res, error)
- func ExponentialJitterOr(base, maxDelay time.Duration) func(attempt int) time.Duration
- func GateRules(pre *directives.Preprocessed) []directives.GateRule
- func GuardCost(reserver cost.Reserver, estimateMicros int64) action.AnyHook
- func HasStreamAtoms(dsl string, registry *Registry) bool
- func MaxTokensFromCtx(ctx context.Context, def int) int
- func NewExecuteAction(comp *Compiler) *action.BuiltAction[GraphExecReq, GraphExecRes]
- func RegisterBoundary(spec BoundarySpec)
- func RegisterPipelines(reg *action.Registry, pipelines []Pipeline) (*action.Registry, error)
- func RegisterResolver(r ServiceResolver)
- func ResetResolvers() int
- func ResolveParamRef(raw string, opts *compileOptions) any
- func SanitizeDSL(rawContent string) string
- func SetDefaultRegistry(r *Registry)
- func StandardLibrary() action.Library
- func WithApprovalGate(g ApprovalGate) func(*Compiler)
- func WithGates(ctx context.Context, gates []directives.GateRule) context.Context
- func WithHooks(hooks ...action.AnyHook) func(*Compiler)
- func WithJournal(j journal.BranchJournal) func(*Compiler)
- func WithMaxTokens(ctx context.Context, n int) context.Context
- func WithReserver(r cost.Reserver) func(*Compiler)
- func WithSecurityObserver(o SecurityObserver) func(*Compiler)
- type ActionMeta
- type ApprovalGate
- type BoundarySpec
- type BranchMode
- type CapabilitySpec
- type CompileOption
- type CompiledEdge
- type CompiledGraph
- type Compiler
- type Config
- type EdgeSpec
- type EffectClass
- type FanInPolicy
- type GraphDefinition
- type GraphExecReq
- type GraphExecRes
- type GraphPolicy
- type Kind
- type Layer
- type Library
- type Metadata
- type NamedOperator
- type NodeKind
- type NodeResult
- type NodeSpec
- type Pipeline
- type Preprocessed
- type Profile
- type ProfilePolicy
- type RecoveryPolicy
- type RecoveryStrategy
- type Registry
- func (r *Registry) Get(name string) (action.AnyAction, bool)
- func (r *Registry) GetOperator(name string) (NamedOperator, bool)
- func (r *Registry) GetStream(name string) (action.AnyStreamAction, bool)
- func (r *Registry) Names() []string
- func (r *Registry) Register(lib Library) error
- func (r *Registry) Resolve(name string) (Kind, bool)
- type Requirement
- type Resolved
- type RetryPolicy
- type Runner
- type SecurityEvent
- type SecurityObserver
- type ServiceResolver
- type SpawnedNode
- type State
- type StreamDrainResult
- type StreamOperator
- type SystemAssembler
- func (s *SystemAssembler) Actions() []action.AnyAction
- func (s *SystemAssembler) AssembleFile(path string) ([]action.AnyAction, error)
- func (s *SystemAssembler) AssembleManifest(manifestDSL string) ([]action.AnyAction, error)
- func (s *SystemAssembler) Register(name string, act action.AnyAction) *SystemAssembler
- type WorkflowFixture
- type WorkflowResult
Constants ¶
const APIVersion = "nexss.ai/v1"
Variables ¶
var DefaultProfiles = map[Profile]ProfilePolicy{ ProfileTrusted: { MaxRetries: 2, }, ProfileUntrustedInput: { MaxRetries: 2, DefaultHooks: []string{ "guard.prompt_injection", "guard.pii_redact", }, }, ProfileNetworkIsolated: { MaxRetries: 1, DefaultHooks: []string{ "guard.network_ssrf", }, AllowedCapabilities: []string{ "agent.*", "ai.*", "sandbox.*", "log.*", "bench.*", }, }, }
DefaultProfiles is the built-in policy table. It is a package-level variable, not a constant, so downstream code can extend it with custom profiles at init time when a new security envelope is genuinely needed. It is not intended to be mutated at runtime.
Functions ¶
func AcquireStateFromGraphState ¶
AcquireStateFromGraphState converts a graph.State into a pooled dag.State without manual map copying.
func AsCostHook ¶
AsCostHook provides an alias for GuardCost to attach cost governance hooks to actions.
func BaseLibrary ¶ added in v0.8.0
BaseLibrary returns a small set of always-useful actions that any .nflow file can rely on without project-specific Go registration.
It is deliberately separate from StandardLibrary:
- StandardLibrary is orchestration mechanics (log, bench, distribute, dispatch, supervisor).
- BaseLibrary is small data and environment primitives that make a pipeline self-contained — stubs, shape transforms, config from env, generated IDs, dynamic dispatch, and a way to force an error.
Applications should mount both:
libs := []action.Library{
flow.BaseLibrary(),
flow.StandardLibrary(),
myProjectLibrary,
}
func BuildCatalogAction ¶
func BuildResolvedConfig ¶ added in v0.8.0
func BuildResolvedConfig(pre *core.Preprocessed, cliArgs []string) map[string]string
BuildResolvedConfig resolves configuration by applying precedence: CLI flags > Environment variables > DSL @config block/lines > defaults.
func CoerceLiteralValue ¶ added in v0.8.0
CoerceLiteralValue parses JSON arrays, objects, booleans, and numbers into native types.
func CompilePipeline ¶
func CompilePipeline(expr string, reg *action.Registry, opts ...CompileOption) (*action.Builder[any, any], error)
CompilePipeline parses a compact arrow expression and returns a builder that runs it. Options are applied before the AST walk; see compile_options.go for what can be configured.
If a flow registry is attached via WithFlowRegistry, the AST is pre-scanned for stream atoms. Pipelines containing any registered source or operator are compiled by the stream compiler; all others take the existing unary path.
func CompileSaga ¶
CompileSaga parses Arrow DSL with embedded transaction rollbacks into a Saga Node.
func ContributeRegistry ¶ added in v0.8.0
ContributeRegistry lets every directive that implements core.RegistryContributor add actions the pipeline can call by name.
Called by the runner between preprocessing and pipeline compilation. Returns base unchanged when no contributor has anything to add.
func CountStreamResult ¶ added in v0.8.0
CountStreamResult is a convenience helper that extracts the item count from a value produced by a drained stream pipeline. It returns 0 for values that do not carry a Count field.
func ExponentialJitterOr ¶ added in v0.6.0
func GateRules ¶ added in v0.8.0
func GateRules(pre *directives.Preprocessed) []directives.GateRule
GateRules returns the @gate rules that were parsed into the Preprocessed structure. Provided as a helper so callers that hold a Preprocessed but not a compiler can still inspect the rules.
This function does NOT belong in the compiler hot path; it exists only because a couple of tools (runner trace, doctor command) need to iterate gates without re-parsing.
func GuardCost ¶ added in v0.3.0
GuardCost returns an action hook that reserves budget before execution and commits or releases it upon completion based on execution outcome.
func HasStreamAtoms ¶ added in v0.8.0
HasStreamAtoms reports whether dsl references any atom the flow registry knows as a stream source or stream operator.
Stream pipelines look like ordinary atom chains to ParseArrowDSL — they parse into a DAG of nodes with no edges and no loop, no conditional, no fallback. But stream atoms are not unary actions; resolving fs.walk through the DAG compiler fails because fs.walk is a source, not an action. This check lets the executor route stream pipelines through the pipeline compiler before the DAG compiler ever sees them.
A nil or empty registry returns false.
func NewExecuteAction ¶
func NewExecuteAction(comp *Compiler) *action.BuiltAction[GraphExecReq, GraphExecRes]
func RegisterBoundary ¶ added in v0.8.0
func RegisterBoundary(spec BoundarySpec)
RegisterBoundary registers a boundary spec. It panics on duplicate names or on a nil Consume, because these are programmer errors that should surface during init rather than at compile time.
Registration is intended to happen in package init functions and in tests. It is not safe for concurrent use; callers must ensure that registration completes before any pipeline is compiled.
func RegisterPipelines ¶ added in v0.6.0
func RegisterResolver ¶ added in v0.8.0
func RegisterResolver(r ServiceResolver)
RegisterResolver is called from init() of a domain package. Resolvers are tried in registration order.
func ResetResolvers ¶ added in v0.8.0
func ResetResolvers() int
ResetResolvers removes every registered resolver and returns the count that was removed. It exists so that tests can isolate themselves from resolvers registered by other tests in the same binary; production code has no reason to call it.
func ResolveParamRef ¶ added in v0.8.0
ResolveParamRef evaluates @config.key, @env.NAME, @flag.name, @arg.N, or returns literal.
func SanitizeDSL ¶ added in v0.5.0
SanitizeDSL strips comments, directives, and declaration headers from a .nflow source so that only the executable pipeline remains.
func SetDefaultRegistry ¶ added in v0.8.0
func SetDefaultRegistry(r *Registry)
SetDefaultRegistry installs a process-wide registry used by CompilePipeline when no explicit registry is passed via WithFlowRegistry.
Intended to be called once at application boot:
reg := flow.NewRegistry() _ = reg.Register(fsio.Library()) flow.SetDefaultRegistry(reg)
Passing nil clears the default. Safe for concurrent use.
func StandardLibrary ¶ added in v0.6.0
func WithApprovalGate ¶
func WithApprovalGate(g ApprovalGate) func(*Compiler)
WithApprovalGate installs the approval gate. A nil argument is ignored: a graph that declares approval requirements but has no gate fails at compile time with a clear message, not silently here.
func WithGates ¶ added in v0.8.0
WithGates attaches the parsed @gate rules to the compilation context. The runner calls this before Execute; a direct caller of Compiler may skip it and gates will be empty.
An empty slice is a no-op so callers do not need to guard the call.
func WithHooks ¶ added in v0.6.0
WithHooks appends hooks that run on every node in the graph. They are applied after the profile hooks, so a caller-provided hook can observe what a profile hook decided.
func WithJournal ¶
func WithJournal(j journal.BranchJournal) func(*Compiler)
WithJournal installs a branch journal. A nil argument is ignored so callers can pass an optionally-nil dependency without a guard.
func WithReserver ¶ added in v0.3.0
WithReserver installs the cost reserver. Nil is a valid value and disables per-node budget reservation.
func WithSecurityObserver ¶ added in v0.8.0
func WithSecurityObserver(o SecurityObserver) func(*Compiler)
WithSecurityObserver installs an observer for profile hook events. Nil is ignored.
Types ¶
type ActionMeta ¶ added in v0.6.0
type ActionMeta = directives.ActionMeta
type ApprovalGate ¶
type ApprovalGate interface {
Check(ctx context.Context, actionName, argsJSON, token string) error
}
ApprovalGate is the interface the compiler uses to ask for human approval before running a node marked as high-risk. It matches runner.TerminalApprovalGate and the mock gate used in tests.
type BoundarySpec ¶ added in v0.8.0
type BoundarySpec struct {
// Name is the atom name recognised in a pipeline, e.g. "collect".
Name string
// Description is a short human-readable summary used by diagnostics
// and documentation.
Description string
// Consume wraps a stream source into a unary action. The returned
// builder produces exactly one value per invocation; that value is
// forwarded to the next node in the pipeline as if it were the
// request of an ordinary unary action.
//
// Consume must not mutate the upstream source. It may wrap the
// source in an adapter that drains the stream once per request.
Consume func(upstream action.AnyStreamAction) *action.Builder[any, any]
}
BoundarySpec describes a stream-to-unary adapter. Boundary atoms are the only places where a flow pipeline switches from stream mode to unary mode.
A pipeline is a left-to-right sequence of nodes connected by `->`. When a boundary atom appears, the compiler:
- compiles everything to its left as a stream source,
- asks the boundary to turn that stream source into a unary action,
- compiles everything to its right as an ordinary unary chain, and
- pipes the three pieces together with a plain sequential pipe.
Boundary atoms never appear in the flow registry, never carry modifiers, and never appear in the action catalog. They are pure compiler directives that happen to share the atom surface syntax.
Precedence ¶
If an atom name is both registered in the flow registry and present in the boundary registry, the boundary wins. This is intentional: boundary names are reserved by convention, and using a boundary name for a normal action is a mistake.
func BoundaryByName ¶ added in v0.8.0
func BoundaryByName(name string) (BoundarySpec, bool)
BoundaryByName looks up a boundary by its atom name. The second return value reports whether a boundary with that name exists.
func NamedBoundaries ¶ added in v0.8.0
func NamedBoundaries() []BoundarySpec
NamedBoundaries returns every registered boundary sorted by name. The result is a copy; the caller may modify it freely.
type BranchMode ¶
type BranchMode string
const ( BranchFirstMatch BranchMode = "first_match" BranchAllMatches BranchMode = "all_matches" )
type CapabilitySpec ¶
type CapabilitySpec struct {
Name string `json:"name"`
Description string `json:"description"`
Route string `json:"route,omitempty"`
Method string `json:"method,omitempty"`
InputSchema map[string]any `json:"input_schema"`
OutputSchema map[string]any `json:"output_schema"`
Tags []string `json:"tags,omitempty"`
IsSystem bool `json:"is_system"`
}
func ExtractCapabilities ¶
func ExtractCapabilities(registry *action.Registry) []CapabilitySpec
type CompileOption ¶ added in v0.8.0
type CompileOption func(*compileOptions)
func WithAtomAdvisors ¶ added in v0.8.0
func WithAtomAdvisors(pre *core.Preprocessed) CompileOption
WithAtomAdvisors collects advisors from every registered AtomPolicy and attaches them to the compilation. Applications that preprocess a file themselves and then call CompilePipeline directly can pass the result here.
func WithCLIArgs ¶ added in v0.8.0
func WithCLIArgs(args []string) CompileOption
func WithConfig ¶ added in v0.8.0
func WithConfig(cfg map[string]string) CompileOption
func WithFlowRegistry ¶ added in v0.8.0
func WithFlowRegistry(reg *Registry) CompileOption
func WithSchemas ¶ added in v0.8.0
func WithSchemas(pre *core.Preprocessed) CompileOption
type CompiledEdge ¶
type CompiledGraph ¶
type CompiledGraph struct {
Definition GraphDefinition
Layers [][]string
NodeByID map[string]NodeSpec
Edges []CompiledEdge
Outgoing map[string][]CompiledEdge
EdgeKeys []string
}
func Compile ¶
func Compile(def GraphDefinition) (*CompiledGraph, error)
func LoadYAML ¶
func LoadYAML(data []byte) (*CompiledGraph, error)
func LoadYAMLFile ¶
func LoadYAMLFile(path string) (*CompiledGraph, error)
func (*CompiledGraph) SelectOutgoing ¶
func (g *CompiledGraph) SelectOutgoing( source string, matches func(condition string) (bool, error), ) ([]CompiledEdge, error)
func (*CompiledGraph) SelectOutgoingDurable ¶
func (g *CompiledGraph) SelectOutgoingDurable( ctx context.Context, j journal.BranchJournal, runID, source string, state *State, ) ([]CompiledEdge, error)
type Compiler ¶
type Compiler struct {
// contains filtered or unexported fields
}
Compiler turns a GraphDefinition into a runnable dag.DAG.
A Compiler is safe for concurrent use by multiple goroutines: every field is read-only after NewCompiler returns, and per-call state lives in the ctx and in the returned CompiledGraph.
func NewCompiler ¶
NewCompiler builds a Compiler with sensible defaults: an in-memory branch journal, no approval gate, no cost reserver, no extra hooks, no security observer. Every default can be overridden with an option.
func (*Compiler) Compile ¶
func (c *Compiler) Compile(ctx context.Context, def GraphDefinition) (*dag.DAG, *CompiledGraph, error)
Compile turns a GraphDefinition into a runnable DAG.
The pipeline is:
- Merge @gate rules from the compilation context into the policy (applyGates). This must happen before Compile(def) freezes the definition into CompiledGraph.Definition.
- Compile the definition into a CompiledGraph (topological sorting, edge validation, cycle detection).
- Resolve the profile and its hooks.
- Enforce the capability allowlist per node.
- Enforce "approval required but no gate configured".
- For every node, build a dag action that: - respects the per-node concurrency limiter, - honours incoming gate outputs, - assembles the request from Params, InputBindings, and upstream node outputs, - enforces MaxContextBytes, - checks the approval gate when required, - places the registry and compiler into the execution context, - decodes the request into the action's typed shape.
- Apply timeout, retry, cost reservation, profile hooks, and compiler-wide hooks in that order.
- For every conditional edge, insert a gate action that picks the outgoing edges using the journal (durable) or the evaluator.
The returned CompiledGraph is what callers inspect for topology, layers, and per-node metadata.
func (*Compiler) CompilePipeline ¶ added in v0.7.0
func (c *Compiler) CompilePipeline(expr string) (action.Executable, error)
CompilePipeline compiles a single expression string against the compiler's registry. Used by nodes that need to compile a child pipeline at run time (supervisor, dispatch fallback).
type Config ¶ added in v0.6.0
type Config struct {
Verbosity int
BudgetMicros int64
Approval string // "danger" | "all" | "none" — validated by the runner
MaxTokens int // 0 = per-node default
Observe string // "live" | "json" | "off"
Provider string
Model string
Sandbox string
OutputFormat string // "" | "json" | "text"
OutputDir string
}
Config is the complete set of runtime knobs for a flow run.
Zero values are meaningful defaults supplied by defaultConfig. Every field is optional in the .nflow, in the environment, and on the CLI.
type EffectClass ¶
type EffectClass string
const ( EffectReadOnly EffectClass = "read_only" EffectSideEffect EffectClass = "side_effect" EffectHighRisk EffectClass = "high_risk" )
type FanInPolicy ¶
type GraphDefinition ¶
type GraphDefinition struct {
APIVersion string `json:"apiVersion" yaml:"apiVersion"`
Kind string `json:"kind" yaml:"kind"`
Metadata Metadata `json:"metadata" yaml:"metadata"`
Profile Profile `json:"profile,omitempty" yaml:"profile,omitempty"`
Policy GraphPolicy `json:"policy,omitempty" yaml:"policy,omitempty"`
Nodes []NodeSpec `json:"nodes" yaml:"nodes"`
Edges []EdgeSpec `json:"edges" yaml:"edges"`
}
func ParseArrowDSL ¶
func ParseArrowDSL(name, dsl string) (GraphDefinition, error)
ParseArrowDSL converts an Arrow DSL expression into a GraphDefinition. It runs the compiler's parser, then walks the resulting AST and produces one NodeSpec per atom, one EdgeSpec per arrow.
Modifiers attached to an atom are parsed a second time here (the parser already collected them as raw strings) so that a small, well-defined subset can be lifted onto NodeSpec fields: effect, approval, timeout, retry. Everything else stays in the atom and is applied at compile time by resolveDynamicNode.
type GraphExecReq ¶
type GraphExecReq struct {
DSL string `json:"dsl,omitempty" cli:"dsl,d" usage:"Compact arrow pipeline expression"`
YAML string `json:"yaml,omitempty" cli:"yaml,y" usage:"Declarative YAML graph manifest"`
Profile Profile `json:"profile,omitempty" usage:"Security envelope: trusted | untrusted_input | network_isolated"`
Name string `json:"name,omitempty" usage:"Source identifier for error messages (usually the .flow path)"`
InitialPayload map[string]any `json:"initial_payload,omitempty" usage:"Initial state values passed to root nodes"`
// Preprocessed carries every declaration the runner parsed from
// the source file. Extension points (atom advisors, pipeline
// wrappers) read their own declarations from it.
Preprocessed *core.Preprocessed `json:"-"`
}
type GraphExecRes ¶
type GraphPolicy ¶
type GraphPolicy struct {
MaxParallelNodes int `json:"max_parallel_nodes,omitempty" yaml:"max_parallel_nodes,omitempty"`
MaxContextBytes int64 `json:"max_context_bytes,omitempty" yaml:"max_context_bytes,omitempty"`
BudgetMicros int64 `json:"budget_micros,omitempty" yaml:"budget_micros,omitempty"`
ApprovalRequiredFor []EffectClass `json:"approval_required_for,omitempty" yaml:"approval_required_for,omitempty"`
FanInRecovery RecoveryPolicy `json:"fan_in_recovery,omitempty" yaml:"fan_in_recovery,omitempty"`
}
type Kind ¶ added in v0.8.0
type Kind uint8
Kind identifies the role a registered entry plays in a pipeline.
type Layer ¶ added in v0.6.0
type Layer uint8
Layer identifies where a resolved config value came from. The runner prints the layer next to each value at -vvv so the operator can see whether a knob came from the CLI, the environment, the .nflow file, or a built-in default.
type Library ¶ added in v0.6.0
Alias types to guarantee complete interoperability with the kernel without duplication.
type NamedOperator ¶ added in v0.8.0
type NamedOperator = action.NamedOperator
Alias types to guarantee complete interoperability with the kernel without duplication.
type NodeResult ¶
type NodeResult[Res any] struct { Spawned SpawnedNode `json:"spawned"` Value Res `json:"value"` Err error `json:"error,omitempty"` }
type NodeSpec ¶
type NodeSpec struct {
ID string `json:"id" yaml:"id"`
Kind NodeKind `json:"kind" yaml:"kind"`
Capability string `json:"capability" yaml:"capability"`
Params map[string]any `json:"params,omitempty" yaml:"params,omitempty"`
InputBindings map[string]string `json:"inputs,omitempty" yaml:"inputs,omitempty"`
InputSchema string `json:"input_schema,omitempty" yaml:"input_schema,omitempty"`
OutputSchema string `json:"output_schema,omitempty" yaml:"output_schema,omitempty"`
Prompt string `json:"prompt,omitempty" yaml:"prompt,omitempty"`
Retry RetryPolicy `json:"retry,omitempty" yaml:"retry,omitempty"`
TimeoutMS int64 `json:"timeout_ms,omitempty" yaml:"timeout_ms,omitempty"`
MaxAttempts int `json:"max_attempts,omitempty" yaml:"max_attempts,omitempty"`
EstimateMicros int64 `json:"estimate_micros,omitempty" yaml:"estimate_micros,omitempty"`
Effect EffectClass `json:"effect,omitempty" yaml:"effect,omitempty"`
Approval bool `json:"approval_required,omitempty" yaml:"approval_required,omitempty"`
BranchMode BranchMode `json:"branch_mode,omitempty" yaml:"branch_mode,omitempty"`
}
type Pipeline ¶ added in v0.6.0
type Pipeline = directives.Pipeline
type Preprocessed ¶ added in v0.6.0
type Preprocessed = directives.Preprocessed
Preprocessed is the alias for directives.Preprocessed. Callers that import flow see flow.Preprocessed; the implementation lives in flow/directives so the directive registry and the data shape stay together.
func Preprocess ¶ added in v0.6.0
func Preprocess(path string) (*Preprocessed, error)
Preprocess reads a .nflow file from disk and resolves every directive: @profile, @include, @pipeline … @end, @action, @description, @require. The resulting DSL is line-aligned with the source file so parser errors point at the original line numbers.
func PreprocessBytes ¶ added in v0.8.0
func PreprocessBytes(source []byte, name string) (*Preprocessed, error)
PreprocessBytes processes an in-memory .nflow source. It is the entry point for embedded libraries (go:embed) that carry their pipeline as a string rather than as a file on disk.
@include is rejected in this mode: an embedded source has no on-disk base directory, so relative paths cannot be resolved.
type Profile ¶ added in v0.8.0
type Profile string
Profile is a declarative security envelope applied to an entire flow at compile time. It is the single place where a user tells the runtime "this flow handles untrusted input" or "this flow talks to the network" without having to wire guards manually.
Profiles are the mechanism behind the "guards are automatic" invariant: the user writes @profile:untrusted_input at the top of the file, and every node in the file receives the profile's hooks, retry policy, and capability restrictions.
The zero value ("" or "trusted") is the permissive default. Every non-default profile is opt-in and documented.
const ( // ProfileTrusted is the zero-value, permissive default. Use for // internal tooling, local scripts, and flows that never touch // untrusted input or network. ProfileTrusted Profile = "trusted" // ProfileUntrustedInput is for flows whose input comes from a // user, an HTTP body, a file, or an LLM. The compiler injects // prompt-injection filtering and PII redaction hooks on every // node that touches LLM traffic. ProfileUntrustedInput Profile = "untrusted_input" // ProfileNetworkIsolated is for flows that must not reach the // public internet. The compiler injects an anti-SSRF hook and, // when the flow also declares untrusted input, refuses to compile // if any node uses :remote= or :exec=. ProfileNetworkIsolated Profile = "network_isolated" )
type ProfilePolicy ¶ added in v0.8.0
type ProfilePolicy struct {
// MaxRetries is the automatic retry budget for transient and
// validation failures. Zero disables automatic retry; the user
// can still write :retry=N to override per-node.
MaxRetries int
// DefaultHooks lists action names that the compiler resolves
// from the registry and attaches as Before-hooks on every node
// in the flow. Missing hooks are a compile-time error: a flow
// that asks for untrusted_input isolation but does not provide
// the guard actions cannot be compiled.
DefaultHooks []string
// AllowedCapabilities is a list of glob patterns matched against
// node capabilities. An empty list means "everything allowed",
// which is the trusted default. A non-empty list means "only
// matching capabilities may appear in this flow".
//
// Patterns are simple globs: "agent.*" matches agent.planner,
// agent.architect, etc. "*" matches everything.
AllowedCapabilities []string
}
ProfilePolicy is the concrete, machine-readable consequence of a Profile. The compiler reads it once, at Compile time, and applies every field to every node in the flow.
This struct is deliberately data-only: no functions, no interfaces, no side effects. It is JSON-serializable so it can be inspected in -vvv trace, unit-tested as a fixture, and eventually loaded from a user-supplied YAML.
func LookupProfile ¶ added in v0.8.0
func LookupProfile(name string) (ProfilePolicy, error)
LookupProfile resolves a profile name to its policy. An empty name resolves to ProfileTrusted. Unknown names are a hard error: a typo in @profile: must not silently fall back to the permissive default.
func (ProfilePolicy) Allows ¶ added in v0.8.0
func (p ProfilePolicy) Allows(capability string) bool
Allows reports whether capability is permitted by this policy. Empty AllowedCapabilities means "everything allowed". Patterns are simple globs matched with a trailing-* shortcut, which covers the "agent.*" / "*" cases users actually write.
type RecoveryPolicy ¶
type RecoveryPolicy struct {
Strategy RecoveryStrategy `json:"strategy,omitempty" yaml:"strategy,omitempty"`
MaxAttempts int `json:"max_attempts,omitempty" yaml:"max_attempts,omitempty"`
BackoffMS int64 `json:"backoff_ms,omitempty" yaml:"backoff_ms,omitempty"`
MaxBackoffMS int64 `json:"max_backoff_ms,omitempty" yaml:"max_backoff_ms,omitempty"`
RetryTransientOnly bool `json:"retry_transient_only,omitempty" yaml:"retry_transient_only,omitempty"`
}
type RecoveryStrategy ¶
type RecoveryStrategy string
const ( RecoveryFailFast RecoveryStrategy = "fail_fast" RecoveryRetryFailed RecoveryStrategy = "retry_failed" RecoveryContinuePartial RecoveryStrategy = "continue_partial" )
type Registry ¶
type Registry struct {
// contains filtered or unexported fields
}
Registry holds registered entries by name.
func DefaultRegistry ¶ added in v0.8.0
func DefaultRegistry() *Registry
DefaultRegistry returns the process-wide registry, if any. Returns nil when unset.
func NewRegistry ¶
func NewRegistry() *Registry
func RegistryFromActionRegistry ¶ added in v0.8.0
RegistryFromActionRegistry bridges a kernel action.Registry into a flow.Registry, copying actions, stream sources, and operators so stream pipelines resolve seamlessly.
func (*Registry) GetOperator ¶ added in v0.8.0
func (r *Registry) GetOperator(name string) (NamedOperator, bool)
type Requirement ¶ added in v0.6.0
type Requirement = directives.Requirement
type Resolved ¶ added in v0.6.0
Resolved pairs the final Config with per-field provenance.
func ResolveConfig ¶ added in v0.6.0
ResolveConfig applies CLI > env > DSL > default to every knob.
cliArgs is the raw flag slice; the parser ignores anything that is not a config knob, so it is safe to pass the same args the runner also uses for runner-specific flags (-i, --assert=, …).
type RetryPolicy ¶
type Runner ¶ added in v0.6.0
type Runner interface {
// RunFlow executes the flow file at path.
//
// args is the raw CLI flag slice; config knobs are parsed
// downstream by flow.ResolveConfig.
// libs are the libraries whose actions the flow may call.
// stdout receives the human-readable trace and metrics table.
// stderr receives fatal errors and warnings.
//
// Returns a process exit code: 0 on success, non-zero otherwise.
RunFlow(
ctx context.Context,
path string,
payload map[string]any,
args []string,
libs []action.Library,
stdout, stderr io.Writer,
) int
}
Runner is the shape of a flow execution engine.
The default implementation is flow/runner.Default. Products that need a different engine — remote execution, a custom observer, a custom approval flow — implement this interface and swap it in.
The interface is deliberately minimal: it captures only what callers actually need (flow path, initial payload, raw CLI args, libraries, and where to write output). The implementation owns the registry, the observer, the approval gate, and every other detail.
Callers that only need the default runner can import github.com/nexssp/flow/runner directly and skip this interface.
type SecurityEvent ¶ added in v0.8.0
type SecurityEvent struct {
Hook string // "guard.prompt_injection"
Node string // "agent.architect"
Phase string // "before" | "after"
Elapsed time.Duration // duration of the hook call itself
Err error // non-nil when the hook rejected the request
}
SecurityEvent is one observation of a profile hook firing around a node. Every profile hook (guard.prompt_injection, guard.pii_redact, guard.network_ssrf, …) emits one "before" and one "after" event per node it wraps. A guard that rejects produces "before" with Err set.
The event is deliberately small and flat: it is meant to be rendered as a single log line, stored as a JSONL row, or counted as a metric without further processing.
type SecurityObserver ¶ added in v0.8.0
type SecurityObserver interface {
OnSecurity(ctx context.Context, ev SecurityEvent)
}
SecurityObserver receives SecurityEvents from the compiler. The interface is a single method so any consumer — a runner observer, a test double, an OpenTelemetry bridge — can implement it without touching the flow package.
type ServiceResolver ¶ added in v0.8.0
type ServiceResolver interface {
Match(actionName string, mods dslparse.Modifiers) bool
Resolve(actionName string, mods dslparse.Modifiers, act action.AnyAction) (action.AnyAction, error)
}
ServiceResolver lets a domain package translate action modifiers into a concrete action implementation without flow knowing anything about the domain.
type SpawnedNode ¶
type SpawnedNode struct {
RunID string `json:"run_id"`
SourceNode string `json:"source_node"`
Edge CompiledEdge `json:"edge"`
TargetNode string `json:"target_node"`
Input *State `json:"input"`
SpawnIndex int `json:"spawn_index"`
}
func SpawnSelected ¶
func SpawnSelected(runID, source string, selected []CompiledEdge, input *State) ([]SpawnedNode, error)
SpawnSelected creates durable invocation units for each selected edge.
type State ¶
type State struct {
// contains filtered or unexported fields
}
func FanIn ¶
func FanIn[Res any]( ctx context.Context, input *State, results []NodeResult[Res], policy FanInPolicy, reduce func(context.Context, *State, []NodeResult[Res]) (*State, error), ) (*State, error)
FanIn deterministically orders parallel child outputs by SpawnIndex and calls reduce.
func NewStateFromDAG ¶
type StreamDrainResult ¶ added in v0.8.0
type StreamDrainResult struct {
Count int `json:"count"`
}
StreamDrainResult is the value produced by a stream pipeline that was drained to a single summary without a boundary.
A pipeline like:
fs.walk -> fs.read
(without a `collect` boundary and without any downstream unary node) does not return the individual items. Instead it runs the stream to completion and returns the number of items that were emitted. This is what a caller gets from `Do` when they compile a stream pipeline without a boundary.
The type is exported so that external test packages and consumers can assert on the concrete value without relying on anonymous structs.
type StreamOperator ¶ added in v0.8.0
type StreamOperator = action.StreamOperator
Alias types to guarantee complete interoperability with the kernel without duplication.
func ComposeOperators ¶ added in v0.8.0
func ComposeOperators(ops ...StreamOperator) (StreamOperator, error)
ComposeOperators chains multiple operators into one.
Does NOT validate element types — validation is the compiler's job (before Build). This function assumes the operators already match.
The composed operator runs Apply sequentially.
func NewTypedStreamOperator ¶ added in v0.8.0
func NewTypedStreamOperator[In, Out any]( name string, op action.StreamOp[In, Out], ) StreamOperator
NewTypedStreamOperator wraps a typed operator. Called from NamedOperator.Build, where In/Out are statically known (srcpack code).
type SystemAssembler ¶ added in v0.4.0
type SystemAssembler struct {
// contains filtered or unexported fields
}
func NewAssembler ¶ added in v0.4.0
func NewAssembler(capabilities ...action.AnyAction) *SystemAssembler
func (*SystemAssembler) Actions ¶ added in v0.4.0
func (s *SystemAssembler) Actions() []action.AnyAction
func (*SystemAssembler) AssembleFile ¶ added in v0.4.0
func (s *SystemAssembler) AssembleFile(path string) ([]action.AnyAction, error)
func (*SystemAssembler) AssembleManifest ¶ added in v0.4.0
func (s *SystemAssembler) AssembleManifest(manifestDSL string) ([]action.AnyAction, error)
func (*SystemAssembler) Register ¶ added in v0.4.0
func (s *SystemAssembler) Register(name string, act action.AnyAction) *SystemAssembler
type WorkflowFixture ¶ added in v0.4.0
type WorkflowFixture struct {
// contains filtered or unexported fields
}
func NewWorkflowTest ¶ added in v0.4.0
func (*WorkflowFixture) Execute ¶ added in v0.4.0
func (wf *WorkflowFixture) Execute(input any) *WorkflowResult
func (*WorkflowFixture) WithTimeout ¶ added in v0.4.0
func (wf *WorkflowFixture) WithTimeout(d time.Duration) *WorkflowFixture
type WorkflowResult ¶ added in v0.4.0
type WorkflowResult struct {
// contains filtered or unexported fields
}
func (*WorkflowResult) Duration ¶ added in v0.4.0
func (r *WorkflowResult) Duration() time.Duration
func (*WorkflowResult) ExpectError ¶ added in v0.4.0
func (r *WorkflowResult) ExpectError() *WorkflowResult
func (*WorkflowResult) ExpectSuccess ¶ added in v0.4.0
func (r *WorkflowResult) ExpectSuccess() *WorkflowResult
func (*WorkflowResult) Output ¶ added in v0.4.0
func (r *WorkflowResult) Output() any
Source Files
¶
- action.go
- arrow_dsl.go
- assembler.go
- atom_advisor.go
- base_library.go
- boundary.go
- catalog.go
- compile_options.go
- compile_options_schema.go
- compiler.go
- compiler_ctx.go
- compiler_hooks.go
- compiler_resolve.go
- compiler_segment.go
- compiler_stream.go
- config.go
- config_cli.go
- config_dsl.go
- config_env.go
- cost.go
- default_registry.go
- definition.go
- dsl.go
- dsl_dynamic.go
- dsl_saga.go
- dsl_stream_detect.go
- journal_helpers.go
- library.go
- parallel.go
- pipeline.go
- preprocess.go
- profile.go
- registry.go
- registry_contrib.go
- resolve_params.go
- runner.go
- sanitize.go
- security.go
- service.go
- state.go
- state_adapter.go
- stream_adapter.go
- stream_pipeline.go
- workflow_testkit.go
Directories
¶
| Path | Synopsis |
|---|---|
|
Package directives holds the DSL directive registry and the shipped directives themselves.
|
Package directives holds the DSL directive registry and the shipped directives themselves. |
|
builtin
Package builtin holds the shipped directives, one subdirectory per directive.
|
Package builtin holds the shipped directives, one subdirectory per directive. |
|
builtin/at_on
file: flow/directives/builtin/at_on/on.go
|
file: flow/directives/builtin/at_on/on.go |
|
core
Package core holds the domain-agnostic parts of the directives system.
|
Package core holds the domain-agnostic parts of the directives system. |
|
Package coverage holds the coverage test suite for the flow DSL.
|
Package coverage holds the coverage test suite for the flow DSL. |
|
examples
|
|
|
08_typed_actions
command
|
|
|
Log nodes.
|
Log nodes. |
|
fsio
flow/nodes/fsio/library.go
|
flow/nodes/fsio/library.go |
|
Package runner provides the shared dynamic-flow execution pipeline used by both `nexssflow` (standalone binary) and `nexssp flow` (subcommand).
|
Package runner provides the shared dynamic-flow execution pipeline used by both `nexssflow` (standalone binary) and `nexssp flow` (subcommand). |
|
bootstrap/console
Package console exposes a small, generic, self-contained web UI for any Nexss binary.
|
Package console exposes a small, generic, self-contained web UI for any Nexss binary. |
|
capability
Package capability resolves flow-node names into runnable action.AnyAction values at flow-execution time.
|
Package capability resolves flow-node names into runnable action.AnyAction values at flow-execution time. |
|
testkit
Package testkit provides flow-runner-specific test helpers.
|
Package testkit provides flow-runner-specific test helpers. |
|
showcase
|
|
|
01_adaptive_rl_router
command
|
|
|
02_self_healing_borg
command
|
|
|
03_evolutionary_optimizer
command
|
|
|
04_grand_showcase
command
|
|
|
flow/testkit/flowtest.go
|
flow/testkit/flowtest.go |
|
Package transport is the sole extension point between the domainless Flow compiler and concrete transport adapters.
|
Package transport is the sole extension point between the domainless Flow compiler and concrete transport adapters. |