Documentation
¶
Overview ¶
Package runner provides the shared dynamic-flow execution pipeline used by both `nexssflow` (standalone binary) and `nexssp flow` (subcommand).
It is deliberately independent of any specific transport: callers supply a flow file path, an initial payload, and optionally a list of assertions; the runner handles DSL sanitization, capability resolution, execution, and assertion evaluation.
The pipeline is:
read manifest -> merge @assert: directives with caller assertions -> sanitize DSL (strip comments, headers, inline @-annotations) -> build static registry (AI bundle + canonical aliases) -> parse :remote / :exec / :wasm bindings from the manifest -> materialize proxy actions, install into registry -> compile + execute the flow -> evaluate assertions against the JSON-encoded result
Capability resolution is done by the child package github.com/nexssp/flow/runner/capability.
Package runner provides the Flow runtime: it compiles, executes, and observes pipelines. It never downloads modules and never generates code at runtime. Build-time validation is a separate concern handled by the code generator in cmd/nexssflow.
Index ¶
- Constants
- func ExecSelf(ctx context.Context, binary string, args []string) int
- func ManifestHash(modules []RequiredModule, opts map[string]map[string]string) string
- func MaterializeForTest(pre *directives.Preprocessed, base *action.Registry) (*action.Registry, error)
- func ParseOptions(pre *directives.Preprocessed, importPath string) map[string]string
- func PrintFlowInfo(ctx context.Context, out io.Writer, req Request, reg *action.Registry) int
- func RegisterMaterializer(m DeclMaterializer)
- func RunAssertions(result any, assertions []string, metrics *RunnerObserver, verbosity int) int
- func RunFlowTest(ctx context.Context, path string, reg *action.Registry, opts TestOptions) int
- func RunWithRegistry(ctx context.Context, req Request, reg *action.Registry, ...) int
- func ValidateOptions(opts map[string]string, allowed ...string) error
- type ApprovalMode
- type Checkpoint
- type CheckpointStore
- type DeclMaterializer
- type Default
- type FieldDoc
- type FileCheckpointStore
- type FlowStep
- type MemoryCheckpointStore
- type MetricRecord
- type ModuleAllowlist
- type ObserverHooks
- type Registrar
- type Request
- type RequiredModule
- type RequiresPlan
- type RunnerObserver
- func (o *RunnerObserver) AddSpend(micros int64)
- func (o *RunnerObserver) Emit(_ context.Context, ev observe.Event)
- func (o *RunnerObserver) Hook() action.AnyHook
- func (o *RunnerObserver) LogData(ctx context.Context, phase string, payload any)
- func (o *RunnerObserver) LogLifecycle(ctx context.Context, phase, msg string, fields ...any)
- func (o *RunnerObserver) LogTiming(ctx context.Context, phase, msg string, fields ...any)
- func (o *RunnerObserver) OnSecurity(_ context.Context, ev flow.SecurityEvent)
- func (o *RunnerObserver) Out() io.Writer
- func (o *RunnerObserver) PrintSummary(out io.Writer)
- func (o *RunnerObserver) ProviderTrace(kind, provider, model string, in, out int, costMicro int64, dur time.Duration, ...)
- func (o *RunnerObserver) SetVerbosity(n int)
- func (o *RunnerObserver) TotalSpentMicros() int64
- func (o *RunnerObserver) TotalTokens() int
- func (o *RunnerObserver) Verbosity() int
- type TerminalApprovalGate
- type TestOptions
Constants ¶
const ( VerbositySilent = 0 VerbositySummary = 1 VerbosityTiming = 2 VerbosityLifecycle = 3 VerbosityData = 4 )
Verbosity levels. Each level adds a NEW CATEGORY of information, not just "more of the same".
Silent — no diagnostics, only final output Summary — one-line summary at the end (-v) Timing — per-node timing and counters (-vv) Lifecycle — config merge, hook firing, state transitions (-vvv) Data — payload inspection per atom, truncated (-vvvv)
Levels 0-3 are implemented in v2.0. Level 4 is specified now to avoid a breaking change later, implemented in F2.
Variables ¶
This section is empty.
Functions ¶
func ExecSelf ¶
ExecSelf re-executes the runner binary with the given args. The context is threaded through so the child is killed when the parent receives SIGTERM or its context is otherwise canceled.
func ManifestHash ¶ added in v0.8.0
func ManifestHash(modules []RequiredModule, opts map[string]map[string]string) string
ManifestHash returns a stable identifier of the @require list and its options. The code generator stores this hash in requires_gen.go. On every build, if the computed hash differs from the stored one, the generator regenerates the file. If it matches, no work is done.
opts maps importPath to its option map.
func MaterializeForTest ¶ added in v0.8.0
func MaterializeForTest(pre *directives.Preprocessed, base *action.Registry) (*action.Registry, error)
MaterializeForTest exposes materializeFromPreprocessed to integration tests that live outside the runner package but need to verify that a domain materializer produces the expected actions.
The "ForTest" suffix and the package doc make it clear this is not part of the public runtime API: it is not documented in user guides and its signature may change alongside the internal one.
func ParseOptions ¶ added in v0.8.0
func ParseOptions(pre *directives.Preprocessed, importPath string) map[string]string
ParseOptions extracts the { key: "value", ... } block attached to a @require declaration. Returns a copy of the map so callers cannot mutate the parsed representation.
func PrintFlowInfo ¶
func RegisterMaterializer ¶ added in v0.8.0
func RegisterMaterializer(m DeclMaterializer)
RegisterMaterializer is called from init() of a domain package. Panics on duplicate keys — a load-time failure beats a runtime silent override.
func RunAssertions ¶
func RunAssertions(result any, assertions []string, metrics *RunnerObserver, verbosity int) int
RunAssertions evaluates every assertion against the flow result.
The environment exposed to an assertion is built in this order:
- Every named node output, keyed by its final path segment (tasks.greet.output becomes env["greet"] and every field inside it is flattened one level up so { status: "ok" } becomes env["status"]).
- The final result of the flow, whose keys overwrite anything already present. This makes `n == 64` work whether `n` came from the last node or was threaded through the pipeline.
At verbosity 0, silence is success. At verbosity >= 1, every assertion prints one line.
func RunFlowTest ¶ added in v0.8.0
RunFlowTest executes a .nflow file with the given registry in test mode: no live LLM calls, no live network, no side effects beyond what the mock fixtures declare. It returns a process exit code:
0 every assertion passed and output matches expected.json 1 a @assert: failed, output mismatch, or runtime error 2 invalid test setup (missing input, bad JSON, bad mock)
If a conventional testdata/<name>.input.json exists, it is used automatically. Same for expected.json and mock.json.
func RunWithRegistry ¶
Types ¶
type ApprovalMode ¶
type ApprovalMode string
const ( ApprovalNone ApprovalMode = "none" ApprovalDanger ApprovalMode = "danger" ApprovalAll ApprovalMode = "all" )
type Checkpoint ¶
type CheckpointStore ¶
type DeclMaterializer ¶ added in v0.8.0
type DeclMaterializer interface {
// Key is the same key the domain stores its declarations under in
// Preprocessed.Declarations. "llm", "sandbox", ...
Key() string
// Materialize turns the raw declaration value into actions. The
// value's concrete type is the domain's business: the runner only
// guarantees that it is whatever Preprocessed.Declarations[key]
// contained.
Materialize(decls any, base *action.Registry) ([]action.AnyAction, error)
}
DeclMaterializer converts one domain's declarations into actions.
A domain package (ai/llm, ai/sandbox) registers a materializer via init(). When a .nflow file declares something in that domain — for example @llm router { provider: "deepseek", model: "..." } — the materializer is called to turn the declaration into a runnable action. The name in the .nflow file becomes the action's registry name.
The interface lives in flow/runner, not in flow/directives, because it depends on action.AnyAction. flow/directives is parse-only; the runner is the assembly point.
type Default ¶
type Default struct {
// Observer is optional. When nil, a fresh observer is created
// with verbosity resolved from the flow config.
Observer *RunnerObserver
}
Default is the standard flow.Runner implementation. It builds a registry from the supplied libraries, wires the observer, and delegates to RunWithRegistry.
The zero value is safe to use.
func (Default) RunWithRegistry ¶
type FileCheckpointStore ¶
type FileCheckpointStore struct {
// contains filtered or unexported fields
}
func NewFileCheckpointStore ¶
func NewFileCheckpointStore(baseDir string) (*FileCheckpointStore, error)
func (*FileCheckpointStore) Delete ¶
func (s *FileCheckpointStore) Delete(ctx context.Context, runID string) error
func (*FileCheckpointStore) Load ¶
func (s *FileCheckpointStore) Load(ctx context.Context, runID string) (Checkpoint, bool, error)
func (*FileCheckpointStore) Save ¶
func (s *FileCheckpointStore) Save(ctx context.Context, cp Checkpoint) error
type MemoryCheckpointStore ¶
type MemoryCheckpointStore struct {
// contains filtered or unexported fields
}
func NewMemoryCheckpointStore ¶
func NewMemoryCheckpointStore() *MemoryCheckpointStore
func (*MemoryCheckpointStore) Delete ¶
func (s *MemoryCheckpointStore) Delete(ctx context.Context, runID string) error
func (*MemoryCheckpointStore) Load ¶
func (s *MemoryCheckpointStore) Load(ctx context.Context, runID string) (Checkpoint, bool, error)
func (*MemoryCheckpointStore) Save ¶
func (s *MemoryCheckpointStore) Save(ctx context.Context, cp Checkpoint) error
type MetricRecord ¶
type ModuleAllowlist ¶ added in v0.8.0
ModuleAllowlist maps "import@version" to true. The application supplies this at build time. An empty allowlist is a fail-closed policy: no @require is permitted.
type ObserverHooks ¶
type ObserverHooks struct {
// KnownModel reports whether model is in the local price catalog.
KnownModel func(model string) bool
// TokensAndCost extracts (prompt, completion, cost_micros, currency,
// known) from an arbitrary action result.
TokensAndCost func(res any) (int, int, int64, cost.Currency, bool)
// ResultSummary renders a short human-readable summary of a result.
ResultSummary func(res any) string
// PromptFromRequest extracts a display prompt from a request. When
// nil, only maps with "prompt" or "goal" are recognized.
PromptFromRequest func(req any) string
}
ObserverHooks lets a domain layer teach the observer how to render results it understands. Every hook is optional: nil functions fall through to the domainless defaults in observer_extract.go.
type Registrar ¶ added in v0.8.0
Registrar is the entry point every @require'd library must export.
The generated requires_gen.go calls this once per @require with the options declared in the .nflow file. The library returns a fully formed action.Library; the caller assembles the final registry.
The signature takes no *action.Registry because registries are immutable. A library cannot add actions to an existing registry; it returns its own library and lets the caller compose.
type Request ¶
type Request struct {
Path string
Payload map[string]any
Args []string
Verbosity int
Info bool
Assertions []string
Metrics bool
OutFormat string
OutDir string
BenchNode string
BenchRuns int
CacheDir string
ApprovalMode ApprovalMode
Resume string
Store CheckpointStore
Stdout io.Writer
Stderr io.Writer
}
type RequiredModule ¶ added in v0.8.0
RequiredModule is one @require declaration after validation.
type RequiresPlan ¶ added in v0.8.0
type RequiresPlan struct {
Modules []RequiredModule
}
RequiresPlan is the validated output of AnalyzeRequires.
func AnalyzeRequires ¶ added in v0.8.0
func AnalyzeRequires(pre *directives.Preprocessed, allow ModuleAllowlist) (*RequiresPlan, error)
AnalyzeRequires validates every @require declaration against the allowlist. Returns xerr.NotFound for any unlisted module. Never downloads, compiles, or executes anything.
type RunnerObserver ¶
type RunnerObserver struct {
// contains filtered or unexported fields
}
func NewRunnerObserver ¶
func NewRunnerObserver(out io.Writer, verbosity int) *RunnerObserver
func NewRunnerObserverWithHooks ¶
func NewRunnerObserverWithHooks(out io.Writer, verbosity int, hooks ObserverHooks) *RunnerObserver
func (*RunnerObserver) AddSpend ¶
func (o *RunnerObserver) AddSpend(micros int64)
func (*RunnerObserver) Hook ¶
func (o *RunnerObserver) Hook() action.AnyHook
func (*RunnerObserver) LogData ¶ added in v0.8.0
func (o *RunnerObserver) LogData(ctx context.Context, phase string, payload any)
LogData emits payload inspection. Truncates aggressively — max 200 chars per value, max 5 items per slice/map. Never dumps a 50MB blob to the terminal.
func (*RunnerObserver) LogLifecycle ¶ added in v0.8.0
func (o *RunnerObserver) LogLifecycle(ctx context.Context, phase, msg string, fields ...any)
LogLifecycle emits a lifecycle-level trace line. Format is intentionally strace-like so grep works:
[3] config.resolve @config.sig: false → true (CLI --sig) [3] hook.before ui.banner (global) [3] atom.enter fs.walk dirs=["."] [3] atom.exit fs.walk 123 files (45ms)
func (*RunnerObserver) LogTiming ¶ added in v0.8.0
func (o *RunnerObserver) LogTiming(ctx context.Context, phase, msg string, fields ...any)
LogTiming emits per-node timing/counters. Cheaper than Lifecycle, fires on every atom but only when verbosity >= 2.
func (*RunnerObserver) OnSecurity ¶ added in v0.8.0
func (o *RunnerObserver) OnSecurity(_ context.Context, ev flow.SecurityEvent)
OnSecurity implements flow.SecurityObserver. It is called by the compiler for every profile hook that fires around a node. At verbosity >= 1, a passing hook prints one line; a rejecting hook prints a distinct failure line. At verbosity < 1, no output is produced — the events still fire, they are simply not rendered.
func (*RunnerObserver) Out ¶
func (o *RunnerObserver) Out() io.Writer
func (*RunnerObserver) PrintSummary ¶
func (o *RunnerObserver) PrintSummary(out io.Writer)
func (*RunnerObserver) ProviderTrace ¶
func (*RunnerObserver) SetVerbosity ¶
func (o *RunnerObserver) SetVerbosity(n int)
func (*RunnerObserver) TotalSpentMicros ¶
func (o *RunnerObserver) TotalSpentMicros() int64
func (*RunnerObserver) TotalTokens ¶
func (o *RunnerObserver) TotalTokens() int
func (*RunnerObserver) Verbosity ¶
func (o *RunnerObserver) Verbosity() int
type TerminalApprovalGate ¶
type TerminalApprovalGate struct {
// contains filtered or unexported fields
}
func NewApprovalGate ¶
func NewApprovalGate(mode ApprovalMode) *TerminalApprovalGate
type TestOptions ¶ added in v0.8.0
type TestOptions struct {
PayloadPath string
ExpectPath string
MockPath string
Stdout io.Writer
Stderr io.Writer
Verbosity int
}
TestOptions configures one flow test invocation. Paths are optional; the conventional layout is:
<flow_dir>/
my.nflow
testdata/
my.input.json
my.expected.json
my.mock.json
Any path left empty falls back to the conventional location when the file exists, and is skipped otherwise.
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
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. |
|
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. |
|
Package testkit provides flow-runner-specific test helpers.
|
Package testkit provides flow-runner-specific test helpers. |