runner

package
v0.9.0 Latest Latest
Warning

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

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

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

View Source
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

func ExecSelf(ctx context.Context, binary string, args []string) int

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 PrintFlowInfo(ctx context.Context, out io.Writer, req Request, reg *action.Registry) int

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 ResolveRequires

func ResolveRequires(ctx context.Context, reqs []directives.Requirement, stdout, stderr io.Writer) (string, error)

ResolveRequires weryfikuje deklaracje @require i zwraca ścieżkę do skompilowanej binarki.

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:

  1. 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"]).
  2. 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 RunFlow added in v0.9.0

func RunFlow(ctx context.Context, req Request) int

RunFlow to skrótowy punkt wejścia dla CLI uruchamiających pliki .nflow.

func RunFlowTest added in v0.8.0

func RunFlowTest(ctx context.Context, path string, reg *action.Registry, opts TestOptions) int

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

func RunWithRegistry(ctx context.Context, req Request, reg *action.Registry, observer *RunnerObserver) int

func ValidateOptions added in v0.8.0

func ValidateOptions(opts map[string]string, allowed ...string) error

ValidateOptions rejects unknown keys against a fixed allowlist. Every library should call this at the top of its Register to fail fast on typos or outdated options.

Types

type ApprovalMode

type ApprovalMode string
const (
	ApprovalNone   ApprovalMode = "none"
	ApprovalDanger ApprovalMode = "danger"
	ApprovalAll    ApprovalMode = "all"
)

type Checkpoint

type Checkpoint struct {
	RunID       string         `json:"run_id"`
	Flow        string         `json:"flow"`
	FlowHash    string         `json:"flow_hash"`
	Layer       int            `json:"layer"`
	SavedAt     time.Time      `json:"saved_at"`
	SpentMicros int64          `json:"spent_micros"`
	State       map[string]any `json:"state"`
	Suspended   bool           `json:"suspended,omitempty"`
}

type CheckpointStore

type CheckpointStore interface {
	Save(ctx context.Context, cp Checkpoint) error
	Load(ctx context.Context, runID string) (Checkpoint, bool, error)
	Delete(ctx context.Context, runID string) error
}

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) RunFlow

func (d Default) RunFlow(
	ctx context.Context,
	path string,
	payload map[string]any,
	args []string,
	libs []action.Library,
	stdout, stderr io.Writer,
) int

func (Default) RunWithRegistry

func (d Default) RunWithRegistry(
	ctx context.Context,
	req Request,
	reg *action.Registry,
	obs *RunnerObserver,
) int

type FieldDoc

type FieldDoc struct {
	Name       string `json:"name"`
	Key        string `json:"key"`
	Type       string `json:"type"`
	Required   bool   `json:"required"`
	Validation string `json:"validation,omitempty"`
}

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

type FlowStep

type FlowStep struct {
	Name   string
	Prompt string
}

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 (*MemoryCheckpointStore) Save

type MetricRecord

type MetricRecord struct {
	ExecutionID   string
	Action        string
	Duration      time.Duration
	PromptSnippet string
	PromptTokens  int
	CompTokens    int
	CostMicros    int64
	Currency      cost.Currency
	CostKnown     bool
	Success       bool
}

type ModuleAllowlist added in v0.8.0

type ModuleAllowlist map[string]bool

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

type Registrar func(opts map[string]string, sink observe.Sink) (action.Library, error)

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

type RequiredModule struct {
	Import  string
	Version string
	Local   bool
}

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) Emit

func (o *RunnerObserver) Emit(_ context.Context, ev observe.Event)

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 (o *RunnerObserver) ProviderTrace(
	kind, provider, model string,
	in, out int,
	costMicro int64,
	dur time.Duration,
	err error,
)

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

func (*TerminalApprovalGate) Check

func (g *TerminalApprovalGate) Check(ctx context.Context, actionName, argsJSON, token string) error

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.

Directories

Path Synopsis
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.

Jump to

Keyboard shortcuts

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