flow

package module
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: 40 Imported by: 0

README

Nexss Flow / nexssp/flow

Action orchestration engine, Arrow DSL, and deterministic DAG compiler for Go.

flow compiles a compact text pipeline into a typed, executable graph. Every node in the graph is a nexssp/kernel action — a typed Go function, a remote call, a subprocess, a WebAssembly module, or another pipeline. Flow provides the DSL, the compiler, the runtime, the built-in nodes, and the extension points that let you add your own.

  • Module: github.com/nexssp/flow
  • Go: 1.26
  • Built with nexss open source packages: kernel, cost, transport, transportai, validation, testkit

Contents

  1. Install
  2. Quick start
  3. The .nflow file format
  4. Arrow DSL operators
  5. Directives
  6. Action modifiers
  7. Transports
  8. Configuration
  9. Built-in nodes
  10. Library system
  11. Running flows
  12. Cost governance
  13. Checkpointing and resume
  14. Approval and HITL
  15. Journaling and replay
  16. Snapshots
  17. Observability
  18. Telemetry hot path
  19. State and conditions
  20. YAML graph definition
  21. Learning router
  22. Evolutionary optimizer
  23. Extending Flow

Install

go get github.com/nexssp/flow@latest

For the CLI:

go install github.com/nexssp/flow/cmd/nexssflow@latest

Quick start

package main

import (
    "context"
    "fmt"

    "github.com/nexssp/flow"
    "github.com/nexssp/kernel/action"
)

func main() {
    ctx := context.Background()

    fetch := action.New("user.fetch", func(_ context.Context, id int) (map[string]any, error) {
        return map[string]any{"id": id, "name": "Alice", "tier": "pro"}, nil
    }).Build()

    notify := action.New("email.send", func(_ context.Context, req map[string]any) (string, error) {
        return fmt.Sprintf("sent to %v", req["to"]), nil
    }).Build()

    reg, err := action.NewRegistry(action.Of(fetch, notify))
    if err != nil { t.Fatal(err) }

    pipeline, err := flow.CompilePipeline(
        `user.fetch -> { to: .name, subject: "Welcome " + .tier } -> email.send`,
        reg,
    )
    if err != nil {
        panic(err)
    }

    res, err := pipeline.Build().Do(ctx, 42)
    if err != nil {
        panic(err)
    }

    fmt.Println(res) // sent to Alice
}

The .nflow file format

A .nflow file is a plain text pipeline. It may contain directives (lines starting with @), a route declaration header, and the pipeline body.

# pipeline comments start with # or //
@config:budget_usd=1.00
@config:approval=danger
@assert: success == true

users.solve:route="POST /api/users/solve":status=200
  users.validate
  -> users.enrich
  -> ( users.audit & users.notify )
  -> users.finalize

The route declaration header names the flow and binds its HTTP route. It is not part of the pipeline; the sanitizer removes it before the pipeline is compiled.


Arrow DSL operators

Operator Syntax Meaning
Sequential pipe A -> B or A | B Passes A's output to B
Parallel scatter ( A & B & C ) Runs concurrently, gathers into map[string]any
Fallback chain A || B Tries A; on error runs B (FirstSuccess)
Conditional gate ? target Runs target only when gate's output is truthy
Autonomous loop loop( A ) until( cond ) Repeats until cond is true, bounded at 15 turns
Inline projection { key: expr } Reshapes the payload between nodes (expr-lang)

Operator precedence (tightest to loosest):

primary  { }  ( )  atom  loop ... until ...
&        parallel
-> |     pipe
||       fallback
?        conditional
Worked examples
# RAG pipeline
vector.search
-> { context: .results, question: .user_query }
-> llm.generate_answer
-> ui.render
# CI/CD with conditional routing
git.pull
-> make.build
-> go.test
-> ( state.tests_passed == true ? k8s.deploy : slack.alert_failure )
# Cheap fallback chain
( redis.get || postgres.query || legacy.rest_api )
-> { formatted_data: .raw_json }
-> http.respond
# Parallel fan-out with fan-in projection
stripe.charge
-> ( pdf.generate_invoice & warehouse.trigger_robot )
-> { status: "processing", invoice_url: .pdf.generate_invoice.url }
-> email.send_receipt
# Bounded autonomous loop
{ attempts: 0, message: "turn 1" }
-> loop(
    { attempts: attempts + 1, message: "turn " + string(attempts + 1) }
    -> log.info
) until( attempts >= 3 )

Directives

Directives start a line with @. They are processed at preprocess time and do not appear in the compiled pipeline.

Directive Purpose
@config:key=value Override a runtime knob
@assert: expr Testkit assertion evaluated after the flow finishes
@pipeline Name … @end Declare a named subflow, callable by name
@include ./path.nflow Merge pipelines and requires from another file
@action name Expose this file as a callable action named name
@description "text" Human-readable description for the action
@require ./local/path Import a local Go library
@require module vX.Y.Z Import a published Go library

@include is transitive and cycle-safe. @require local paths resolve to the containing module path by walking up to the nearest go.mod.

Named pipelines
@pipeline greet
  { message: "hello, " + name }
  -> log.info
@end

{ name: "world" } -> greet
{ name: "again" } -> greet

A pipeline is registered as an action under its name and can be called from anywhere in the same manifest or from any file that @includes it.

Requiring libraries
@require ./text
@require github.com/acme/text-tools v1.0.0

{ message: "hello, world" }
-> text_tools.uppercase
-> log.info

Each required package must export func Library() flow.Library.


Action modifiers

Modifiers are :key=value pairs appended to an atom. They configure the action at call time. Order is not significant. Boolean flags have no =value:

agent.critic:model="deepseek-flash":timeout=10s:retry=2:cache=5m
-> users.save:validate:idempotent:status=201
-> tools.search:coalesce:dedup
Routing and identity
Modifier Effect
:route="METHOD /path" Bind an HTTP route
:http="METHOD /path" Alias for :route=
:name=identifier Rename the action
:desc="text" / :description="text" Description override
:type=Req->Res Override request/response type names
:scope=public|internal|system Action scope
:status=code HTTP success status
Resilience
Modifier Effect
:timeout=30s Per-call timeout
:retry=N Max retry attempts
:retry_if=predicate Retry predicate (default: transient only)
:backoff=strategy,base:X,max:Y Backoff strategy
:backoff_base=100ms Backoff base (explicit form)
:backoff_max=30s Backoff ceiling (explicit form)
:breaker=failures:N,cooldown:30s Circuit breaker
:breaker_failures=N Failure threshold
:breaker_cooldown=30s Half-open reset delay
:priority=critical|normal|low Load-shedding tier
:concurrency=N Max concurrent in-flight calls
:rate_limit=N/s Token bucket rate limit
:burst=N Rate-limit burst size
Caching and deduplication
Modifier Effect
:cache=5m Read-through cache TTL
:cache_key=... Custom cache key
:coalesce Share in-flight results across concurrent callers
:dedup Same-key callers block until the first completes
:idempotent Register idempotency metadata
:idempotency_header=X-Key Custom idempotency header name
Security and governance
Modifier Effect
:auth Require an authenticated context
:role=name Require a role
:perm=name Require a permission
:feature=flag Require a feature toggle
:budget_micros=N Per-node cost estimate for the reservation guard
:budget=$1.00 Same, in currency units
:audit Emit an audit record on success
:hitl="prompt" Mark for human-in-the-loop approval
:hitl_options=a,b,c Approval option labels
:hitl_trigger=reason Trigger condition description
Lifecycle and deprecation
Modifier Effect
:deprecated Mark deprecated
:since=v1.2.0 Deprecation version
:use=replacement Suggested replacement
:validate Enable request struct validation
:debug Log every invocation
Transport-specific
Modifier Effect
:channel=name SSE channel
:cli_alias=a,b CLI command aliases
:cli_desc="text" CLI help text
:a2a_desc="text" A2A role description
:a2a_example=text A2A usage example

Transports

Every transport modifier accepts a target. Multiple transports can be attached to the same action.

Modifier Target format Purpose
:route= "METHOD /path" HTTP REST
:sse= "/path" Server-Sent Events
:raw= "METHOD /path" Raw HTTP handler
:cli= "command:help" CLI subcommand
:cron= "every 5m" or "* * * * *" Cron schedule
:worker= "every 30s" Background worker
:topic= "topic.name" In-process bus
:a2a= "role" Agent-to-agent
:nats= "subject" NATS pub/sub
:nats_rpc= "subject" NATS request/reply
:nats_kv= "bucket/key" NATS KV get/watch
:nats_durable= "stream/subject/durable[/dlq]" JetStream durable work
:nats_consumer= "stream/subject/durable" Custom JetStream consumer
:nats_obj= "bucket/pattern" NATS Object Store
:nats_svc= "service/version/endpoint/subject" NATS microservice
Capability bindings

A .nflow file can proxy a node to an external target without writing Go:

Modifier Target Behaviour
:remote="http://..." URL JSON POST request/response
:exec="rg --json" shell command JSON on stdin, JSON on stdout
:wasm="./x.wasm" wasm file JSON on stdin, JSON on stdout (wazero)

The first call compiles the WASM module; subsequent calls reuse it.

agent.worker:remote="http://10.0.0.6:9002/work":timeout=45s
tools.ripgrep:exec="rg --json":timeout=10s
skills.lint:wasm="./skills/lint.wasm":timeout=15s

Configuration

Four layers, applied in this order (later overrides earlier):

CLI  >  env  >  @config:  >  defaults
Knobs
Knob Type Default Purpose
verbosity / v int 0 0–3, clamped
budget_micros int 10,000,000 Hard cost ceiling
budget_usd / budget float — Same, in USD
approval string danger danger, all, none
max_tokens int 0 Per-run LLM token ceiling
observe string live live, json, off
provider string — Default LLM provider
model string — Default model
sandbox string — Sandbox driver
out / output_format string — json, text
out_dir / output_dir string .runs Output directory
Environment variables

NEXSS_VERBOSITY, NEXSS_BUDGET_MICROS, NEXSS_BUDGET_USD, NEXSS_APPROVAL, NEXSS_MAX_TOKENS, NEXSS_OBSERVE, NEXSS_PROVIDER, NEXSS_MODEL, NEXSS_SANDBOX, NEXSS_OUT, NEXSS_OUT_DIR.

CLI flags
-v | -vv | -vvv         verbosity
-q | --quiet            silence
--budget=<usd>          budget in USD
--budget-micros=<n>     budget in micros
--max-tokens=<n>        LLM token ceiling
--approval=<mode>       danger | all | none
--observe=<mode>        live | json | off
--provider=<name>
--model=<name>
--sandbox=<name>
--out=<format>
--out-dir=<path>
-i | --info             describe flow, do not execute
--assert="expr"         testkit assertion
--resume=<runID>        resume from checkpoint
--bench=<node>          benchmark a node
--bench-runs=<n>        benchmark iterations
--cache=<dir>           cache directory

Built-in nodes

flow.StandardLibrary() provides:

Node Purpose
log.info / log.warn / log.error Structured log, passes payload through
bench.run Run another action N times, report latency distribution
bench.save Write a benchmark result to a file
bench.compare Compare current benchmark against a baseline
distribute.map Invoke an action once per item, bounded concurrency
distribute.reduce Fold a distribute.map result
supervisor Compile and run child pipelines on the fly

bench.run, distribute.map, and supervisor resolve their target through the registry the compiler places on the execution context.

Aliases

log, info, warn, error, bench, benchmark, compare, diff, map, fanout, parallel, reduce, fold.

Canonical names always win over aliases; user-provided actions always win over built-in aliases.


Library system

A library is a named bag of actions, hooks, and aliases:

type Library struct {
    Name        string
    Description string
    Actions     []action.AnyAction
    Hooks       []action.AnyHook
    Aliases     []Alias           // Alias{Canonical, Short []string}
    Overrides   []string          // canonical names this library intentionally replaces
}

flow.BuildRegistry(libs...) applies four rules:

  1. Every primary action registers under its canonical name.
  2. When two libraries declare the same canonical name, the later library must list it in Overrides, otherwise BuildRegistry returns an error.
  3. Hooks from every library are applied to every surviving action.
  4. Aliases are registered last, so canonical names always win.
reg, err := flow.BuildRegistry(
    flow.StandardLibrary(),
    myLibrary,
)
Declaring your own library
package mylib

import (
    "github.com/nexssp/flow"
    "github.com/nexssp/kernel/action"
)

func Actions() []action.AnyAction {
    return []action.AnyAction{
        action.New("mylib.echo", func(_ context.Context, s string) (string, error) {
            return s, nil
        }).Tag("mylib").Build(),
    }
}

func Library() flow.Library {
    return flow.Library{
        Name:        "mylib",
        Description: "Example library",
        Actions:     Actions(),
        Aliases: []flow.Alias{
            {Canonical: "mylib.echo", Short: []string{"echo"}},
        },
    }
}

Consume it in a flow file with @require ./mylib, or programmatically:

reg, err := flow.BuildRegistry(
    flow.StandardLibrary(),
    mylib.Library(),
)

Running flows

CLI
nexssflow ./pipeline.nflow '{"user_id": 42}' -vvv --assert="success == true"
nexssflow ./pipeline.nflow --info
nexssflow ./pipeline.nflow --resume=run_1731000000
Programmatic
req := flowrunner.Request{
    Path:    "./pipeline.nflow",
    Payload: map[string]any{"user_id": 42},
    Args:    []string{"-vv"},
    Stdout:  os.Stdout,
    Stderr:  os.Stderr,
}

exit := flowrunner.Default{}.RunFlow(ctx, req.Path, req.Payload, req.Args,
    []flow.Library{flow.StandardLibrary()}, req.Stdout, req.Stderr)

Or with a custom registry:

exit := flowrunner.RunWithRegistry(ctx, req, reg, observer)
Embedded compilation
compiler := flow.NewCompiler(reg,
    flow.WithApprovalGate(gate),
    flow.WithReserver(ledger),
    flow.WithJournal(journal),
    flow.WithHooks(obsHook, liveHook),
)

execAct := flow.NewExecuteAction(compiler)

res, err := execAct.Do(ctx, flow.GraphExecReq{
    DSL:            pipelineDSL,
    InitialPayload: payload,
})
Describing a flow
exit := flowrunner.PrintFlowInfo(ctx, os.Stdout, req, reg)

Prints the pipeline topography, entry payload shape, and per-node metadata.


Cost governance

ledger := cost.NewLedger(10_000_000, cost.USD)   // $10.00

compiler := flow.NewCompiler(reg, flow.WithReserver(ledger))

Attach to individual actions:

guarded := action.New("ai.complete", handler).
    AnyHook(flow.GuardCost(ledger, 50_000)).   // reserve $0.05
    Build()

The compiler reserves the estimated cost of every node before execution. If the budget is exceeded, the node fails before the handler runs.

The ledger reports per-currency totals; currencies are never summed against each other.

Multi-tenant ledgers
type TenantLedgerRegistry struct {
    mu      sync.RWMutex
    ledgers map[string]*cost.Ledger
}

func TenantCostHook(reg *TenantLedgerRegistry, estimate int64) action.AnyHook {
    return action.AnyHook{
        Before: func(ctx context.Context, _ any, _ *action.Meta) (context.Context, error) {
            tenant, _ := TenantFromContext(ctx)
            ledger, ok := reg.Get(tenant.TenantID)
            if !ok {
                return ctx, xerr.Forbidden("tenant has no ledger")
            }
            reservation, err := ledger.Reserve(ctx, estimate)
            if err != nil {
                return ctx, err
            }
            return context.WithValue(ctx, reservationKey{}, reservation), nil
        },
        After: func(ctx context.Context, _ any, result any, actionErr error, _ *action.Meta) {
            // commit or release based on actionErr
        },
    }
}

Checkpointing and resume

Runner.Request.Store is a 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
}

Built-in implementations: FileCheckpointStore (atomic write, 0600), MemoryCheckpointStore.

On failure the runner writes a checkpoint and prints a resume command:

nexssflow ./pipeline.nflow --resume=run_1731000000

Resume replays the checkpoint's state and skips completed layers. A flow_hash guard rejects a checkpoint if the flow has changed.


Approval and HITL

type ApprovalGate interface {
    Check(ctx context.Context, actionName, argsJSON, token string) error
}

runner.TerminalApprovalGate prompts on stdin. Modes:

  • danger — prompts only for actions with exec, write, delete, rm, drop, truncate, migration, patch, or deploy in the name
  • all — prompts for every action
  • none — no prompts

Programmatic callers pass the token via xctx.WithApprovalToken.

A graph compiled with Approval: true nodes or Policy.ApprovalRequiredFor will refuse to compile without a gate.

Custom gates

Any type with a Check method satisfies ApprovalGate:

type SlackApproval struct {
    channel string
}

func (s *SlackApproval) Check(ctx context.Context, actionName, args, token string) error {
    // send message to Slack, wait for reaction, return nil or error
}

Journaling and replay

journal.BranchJournal records which edge was taken for every conditional node. Replaying a run with the same run_id uses the recorded decisions instead of re-evaluating conditions.

type BranchJournal interface {
    Get(ctx context.Context, runID, sourceNode string) ([]BranchRecord, bool, error)
    Put(ctx context.Context, runID, sourceNode string, records []BranchRecord) error
}

Built-in implementations:

  • MemoryBranchJournal — for tests
  • SQLBranchJournal — SQLite-compatible schema (works with SQLite, Postgres, MySQL)

Enable by passing flow.WithJournal(j) to NewCompiler and setting xctx.WithExecutionID(ctx, runID) on the run.

db, _ := sql.Open("sqlite3", ":memory:")
j := journal.NewSQLBranchJournal(db)
_ = j.EnsureSchema(ctx)

compiler := flow.NewCompiler(reg, flow.WithJournal(j))

Snapshots

journal.SnapshotJournal persists a run's state at every layer boundary so a crashed run can be recovered from the last successful layer.

type Snapshot struct {
    RunID       string
    StepIndex   int
    StateData   map[string]any
    SpentMicros int64
    Timestamp   time.Time
}

FileSnapshotJournal writes each step under <baseDir>/<runID>/step_NNNNN.json and a latest.json pointer. Every write is atomic and durable: temp file, fsync, rename, fsync directory.

j, _ := journal.NewFileSnapshotJournal("./snapshots")
_ = j.Save(ctx, journal.Snapshot{
    RunID:     "run_42",
    StepIndex: 3,
    StateData: map[string]any{"step": 3},
})

snap, found, _ := j.Recover(ctx, "run_42")

Observability

sink := observe.NewPrometheusSink()

act := action.New("user.fetch", handler).
    AnyHook(observe.Hook(sink)).
    Build()

Built-in sinks:

Sink Purpose
observe.NewSlogSink(logger) Structured log/slog output
observe.NewMemorySink(cap) Thread-safe ring buffer of recent events
observe.NewMetricsSink() Aggregate counters by action + kind
observe.NewPrometheusSink() Prometheus exposition text
observe.NewJSONLSink(w, maxBytes) Newline-delimited JSON

Event kinds: executed, error, retry, cache_hit, cache_miss, canceled, panic, coalesced, deduplicated.

Every event carries ExecutionID, TraceID, SpanID, TenantID, UserID, Duration, and the full Request/Response payloads.

Fan-out to multiple sinks
type fanoutSink struct{ sinks []observe.Sink }

func (s *fanoutSink) Emit(ctx context.Context, e observe.Event) {
    for _, sink := range s.sinks {
        sink.Emit(ctx, e)
    }
}

sink := &fanoutSink{sinks: []observe.Sink{
    observe.NewSlogSink(slog.Default()),
    observe.NewMetricsSink(),
    observe.NewPrometheusSink(),
}}

act := action.New("user.fetch", handler).AnyHook(observe.Hook(sink)).Build()

Telemetry hot path

telemetry.HotPathExecutor wraps a node and records one 64-byte ring slot per invocation. The ring buffer is lock-free SPSC, cache-line padded, and allocates nothing per push.

ring := ringbuf.New(8192)
exec := telemetry.NewHotPathExecutor(ring)

out, err := exec.ExecuteNode(ctx, nodeID, invoker, payloadBytes)

// batch drain into a caller-owned slice
var batch [64]ringbuf.Slot
n := ring.BatchDrain(batch[:])

State and conditions

flow.NewState(map[string]any) wraps a run's mutable state. State.Get walks nested paths and struct fields, supporting snake_case JSON tags as well as Go field names.

flow.EvaluateCondition(expr, state) evaluates a graph edge condition:

state.tests_passed == true
state.score >= 0.85
state.error exists
state.status != "pending"

Supported operators: ==, !=, >, >=, <, <=, exists. Literals: "string", 'string', true, false, numbers.


YAML graph definition

An alternative to the Arrow DSL. Same compiler, same runtime.

apiVersion: nexss.ai/v1
kind: Graph
metadata:
  name: review_pipeline
  version: "1.0.0"

policy:
  max_parallel_nodes: 8
  max_context_bytes: 1048576
  budget_micros: 5000000
  approval_required_for: ["high_risk"]
  fan_in_recovery:
    strategy: retry_failed
    max_attempts: 3
    backoff_ms: 200
    max_backoff_ms: 5000
    retry_transient_only: true

nodes:
  - id: fetch
    capability: user.fetch
    kind: tool
    timeout_ms: 5000
    retry: { max_attempts: 2, backoff: exponential }
  - id: review
    capability: user.review
    kind: llm
    estimate_micros: 50000
    effect: read_only

edges:
  - from: fetch
    to: review
    when: 'state.tier == "pro"'
    priority: 1
  - from: fetch
    to: notify
    otherwise: true
    priority: 99

Load with flow.LoadYAML(data) or flow.LoadYAMLFile(path).

  • Node kinds: deterministic, llm, tool, subgraph, human, approval
  • Branch modes: first_match (default), all_matches
  • Fan-in recovery: fail_fast, retry_failed, continue_partial

Learning router

A multi-armed bandit that routes each call to one of N candidate actions, learning from reward signals.

router := learn.NewRouter(learn.RouterConfig{
    Name:         "gateway.router",
    Temperature:  1.8,
    LearningRate: 0.15,
    RewardFn:     myReward,
}, providerFast, providerCheap, providerReliable)

out, err := router.DoAny(ctx, payload)

RewardFn receives the result, error, and duration; returns a float64 reward.

rewardFn := func(res any, err error, d time.Duration) float64 {
    if err != nil {
        return -50.0
    }
    return 20.0 - float64(d.Milliseconds())/10.0
}

The router uses online softmax Q-learning and warms up by visiting each candidate once.


Evolutionary optimizer

Searches for the optimal pipeline topology by mutating a baseline DSL across generations.

best, err := optimizer.Evolve(ctx, baselineDSL, reg, evaluator,
    optimizer.Options{Generations: 4, Population: 6})

Mutations: add :retry=N, wrap two adjacent nodes in ( A & B ), add a fallback with ||.

evaluator := func(ctx context.Context, candidate action.AnyAction) (float64, error) {
    start := time.Now()
    _, err := candidate.DoAny(ctx, myPayload)
    if err != nil {
        return -500.0, nil
    }
    return 100.0 - float64(time.Since(start).Milliseconds()), nil
}

The result is a candidate DSL string and its fitness score.


Extending Flow

Adding a custom node

Any action.AnyAction from nexssp/kernel can appear in a flow. Build one with action.New and register it:

myNode := action.New("my.custom_node", func(ctx context.Context, req MyReq) (MyRes, error) {
	return MyRes{}, nil
}).
	Description("Does something custom").
	Tag("custom").
	Build()

reg := action.MustNewRegistry(action.Of(myNode))

Adding a custom transport

Implement transport.Transport:

type Transport interface {
    fmt.Stringer
    CanHandle(b action.Binding) bool
    Mount(actions []action.AnyAction)
    Do(ctx context.Context, v any) (any, error)
}

Register it with the app builder via WithLoader:

app.WithLoader(func(asm *bootstrap.Assembly) error {
    myTransport := myt.New()
    myTransport.Mount(asm.Actions)
    return nil
})
Adding project-specific actions

When a .nflow file references an action that is not in BaseLibrary or StandardLibrary, mount it from a small main.go in your project:

package main

import (
    "context"
    "os"

    "github.com/nexssp/flow"
    "github.com/nexssp/flow/runner"
    "github.com/nexssp/kernel/action"
)

func main() {
    libs := []action.Library{
        flow.BaseLibrary(),
        flow.StandardLibrary(),
        {
            Name: "myproject",
            Actions: []action.AnyAction{
                action.New("myproject.greet", func(_ context.Context, name string) (string, error) {
                    return "hello, " + name, nil
                }).Build(),
            },
        },
    }

    os.Exit(runner.Default{}.RunFlow(
        context.Background(), os.Args[1], map[string]any{},
        os.Args[2:], libs, os.Stdout, os.Stderr,
    ))
}

That is the intended extension point. The shipped nexssflow binary is deliberately closed; applications compose their own.

Adding a hook

Hooks run before and after every node. Use them for logging, auditing, cost tracking, rate limiting, or anything that wraps the call:

auditHook := action.AnyHook{
    Before: func(ctx context.Context, req any, meta *action.Meta) (context.Context, error) {
        log.Printf("calling %s", meta.Name)
        return ctx, nil
    },
    After: func(ctx context.Context, req, res any, err error, meta *action.Meta) {
        if err != nil {
            log.Printf("%s failed: %v", meta.Name, err)
        }
    },
}

compiler := flow.NewCompiler(reg, flow.WithHooks(auditHook))
Adding a middleware

Kernel middlewares wrap a single action's handler:

wrapped := action.New("my.op", handler).
    Timeout(5 * time.Second).
    Retry(3, action.ExponentialJitter(100*time.Millisecond, 2*time.Second)).
    Cache(1*time.Minute, func(r MyReq) string { return r.Key }).
    Dedup(func(r MyReq) string { return r.Key }).
    RateLimit(100, 200).
    ConcurrencyLimit(16).
    Build()

Every middleware that appears in the DSL is a kernel middleware. To add a new one, add the modifier parser branch in dslparse/ and the corresponding Builder method in kernel/action.

Adding a graph node kind

Node kinds are strings validated by definition.go. To add a new kind (e.g. node_webhook), extend NodeKind and add a case in the compiler's resolveCapability switch.

Adding a sink
type Sink interface {
    Emit(context.Context, Event)
}

Register with the compiler via flow.WithHooks(observe.Hook(mySink)), or attach to individual actions with .AnyHook(observe.Hook(mySink)).

Adding a checkpoint store

Implement CheckpointStore. Use any storage backend: S3, Redis, Postgres, memory.

Adding an approval gate

Implement ApprovalGate. Wire it with flow.WithApprovalGate(myGate).

Adding a journal backend

Implement journal.BranchJournal or journal.SnapshotJournal. SQL, file, memory, or anything else.

Adding a config knob

Extend Config in config.go, add the parser branch in each of config_cli.go, config_env.go, config_dsl.go, and add the key to knobKeys.


See also

  • FLOW.en.md — design rationale and manifesto
  • examples/ — runnable .nflow files and Go examples
  • showcase/ — adaptive router, hot swap, evolutionary optimizer, multi-tenant governance
  • runner/ — the runner package in isolation

Apache License 2.0. Copyright © 2018–2026 Marcin Polak and Contributors.

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

View Source
const APIVersion = "nexss.ai/v1"

Variables

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

func AcquireStateFromGraphState(s *State) *dag.State

AcquireStateFromGraphState converts a graph.State into a pooled dag.State without manual map copying.

func AsCostHook

func AsCostHook(reserver cost.Reserver, estimateMicros int64, _ ...int64) action.AnyHook

AsCostHook provides an alias for GuardCost to attach cost governance hooks to actions.

func BaseLibrary added in v0.8.0

func BaseLibrary() action.Library

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 BuildCatalogAction(registry *action.Registry, bindings ...action.Binding) action.AnyAction

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

func CoerceLiteralValue(trimmed string) any

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

func CompileSaga(expr string, reg *action.Registry) (*action.Builder[any, any], error)

CompileSaga parses Arrow DSL with embedded transaction rollbacks into a Saga Node.

func ContributeRegistry added in v0.8.0

func ContributeRegistry(pre *core.Preprocessed, base *action.Registry) (*action.Registry, error)

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

func CountStreamResult(res any) int

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 EvaluateCondition

func EvaluateCondition(condition string, state *State) (bool, error)

func Execute

func Execute[Req, Res any](ctx context.Context, act *action.BuiltAction[Req, Res], req Req) (Res, error)

func ExponentialJitterOr added in v0.6.0

func ExponentialJitterOr(base, maxDelay time.Duration) func(attempt int) time.Duration

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

func GuardCost(reserver cost.Reserver, estimateMicros int64) action.AnyHook

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

func HasStreamAtoms(dsl string, registry *Registry) bool

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 MaxTokensFromCtx added in v0.6.0

func MaxTokensFromCtx(ctx context.Context, def int) int

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 RegisterPipelines(ctx context.Context, reg *action.Registry, pipelines []Pipeline) (*action.Registry, error)

func ResolveParamRef added in v0.8.0

func ResolveParamRef(raw string, opts *compileOptions) any

ResolveParamRef evaluates @config.key, @env.NAME, @flag.name, @arg.N, or returns literal.

func SanitizeDSL added in v0.5.0

func SanitizeDSL(rawContent string) string

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 StandardLibrary() action.Library

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

func WithGates(ctx context.Context, gates []directives.GateRule) context.Context

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

func WithHooks(hooks ...action.AnyHook) func(*Compiler)

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 WithMaxTokens added in v0.6.0

func WithMaxTokens(ctx context.Context, n int) context.Context

func WithReserver added in v0.3.0

func WithReserver(r cost.Reserver) func(*Compiler)

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:

  1. compiles everything to its left as a stream source,
  2. asks the boundary to turn that stream source into a unary action,
  3. compiles everything to its right as an ordinary unary chain, and
  4. 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 WithCompileContext added in v0.9.0

func WithCompileContext(ctx context.Context) 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 CompiledEdge struct {
	From      string
	To        string
	When      string
	Otherwise bool
	Priority  int
}

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

func NewCompiler(reg *action.Registry, opts ...func(*Compiler)) *Compiler

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:

  1. Merge @gate rules from the compilation context into the policy (applyGates). This must happen before Compile(def) freezes the definition into CompiledGraph.Definition.
  2. Compile the definition into a CompiledGraph (topological sorting, edge validation, cycle detection).
  3. Resolve the profile and its hooks.
  4. Enforce the capability allowlist per node.
  5. Enforce "approval required but no gate configured".
  6. 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.
  7. Apply timeout, retry, cost reservation, profile hooks, and compiler-wide hooks in that order.
  8. 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 EdgeSpec

type EdgeSpec struct {
	From      string `json:"from" yaml:"from"`
	To        string `json:"to" yaml:"to"`
	When      string `json:"when,omitempty" yaml:"when,omitempty"`
	Otherwise bool   `json:"otherwise,omitempty" yaml:"otherwise,omitempty"`
	Priority  int    `json:"priority,omitempty" yaml:"priority,omitempty"`
}

type EffectClass

type EffectClass string
const (
	EffectReadOnly   EffectClass = "read_only"
	EffectSideEffect EffectClass = "side_effect"
	EffectHighRisk   EffectClass = "high_risk"
)

type FanInPolicy

type FanInPolicy struct {
	RequireAll   bool
	AllowPartial bool
	FailOnEmpty  bool
}

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 GraphExecRes struct {
	GraphName  string         `json:"graph_name"`
	Outputs    map[string]any `json:"outputs"`
	Result     any            `json:"result,omitempty"`
	LayersRun  int            `json:"layers_run"`
	DurationMS int64          `json:"duration_ms"`
}

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.

const (
	KindUnary Kind = iota
	KindSource
	KindOperator
)

func (Kind) String added in v0.8.0

func (k Kind) String() string

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.

const (
	LayerDefault Layer = iota
	LayerDSL
	LayerEnv
	LayerCLI
)

func (Layer) String added in v0.6.0

func (l Layer) String() string

type Library added in v0.6.0

type Library = action.Library

Alias types to guarantee complete interoperability with the kernel without duplication.

type Metadata

type Metadata struct {
	Name        string `json:"name" yaml:"name"`
	Version     string `json:"version" yaml:"version"`
	Description string `json:"description,omitempty" yaml:"description,omitempty"`
}

type NamedOperator added in v0.8.0

type NamedOperator = action.NamedOperator

Alias types to guarantee complete interoperability with the kernel without duplication.

type NodeKind

type NodeKind string
const (
	NodeDeterministic NodeKind = "deterministic"
	NodeLLM           NodeKind = "llm"
	NodeTool          NodeKind = "tool"
	NodeSubgraph      NodeKind = "subgraph"
	NodeHuman         NodeKind = "human"
	NodeApproval      NodeKind = "approval"
)

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

func Preprocess added in v0.6.0

func Preprocess(path string) (*Preprocessed, error)

func PreprocessBytes added in v0.8.0

func PreprocessBytes(source []byte, name string) (*Preprocessed, error)

func PreprocessFS added in v0.9.0

func PreprocessFS(fsys fs.FS, path string) (*Preprocessed, error)

PreprocessFS pozwala na przetwarzanie manifestów z wirtualnego systemu plików (np. embed.FS) ze wsparciem dla dyrektywy @include.

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

func RegistryFromActionRegistry(actionRegistry *action.Registry) *Registry

RegistryFromActionRegistry bridges a kernel action.Registry into a flow.Registry, copying actions, stream sources, and operators so stream pipelines resolve seamlessly.

func (*Registry) Get

func (r *Registry) Get(name string) (action.AnyAction, bool)

func (*Registry) GetOperator added in v0.8.0

func (r *Registry) GetOperator(name string) (NamedOperator, bool)

func (*Registry) GetStream added in v0.8.0

func (r *Registry) GetStream(name string) (action.AnyStreamAction, bool)

func (*Registry) Names added in v0.8.0

func (r *Registry) Names() []string

func (*Registry) Register added in v0.8.0

func (r *Registry) Register(lib Library) error

func (*Registry) Resolve added in v0.8.0

func (r *Registry) Resolve(name string) (Kind, bool)

type Requirement added in v0.6.0

type Requirement = directives.Requirement

type Resolved added in v0.6.0

type Resolved struct {
	Config     Config
	Provenance map[string]Layer
}

Resolved pairs the final Config with per-field provenance.

func ResolveConfig added in v0.6.0

func ResolveConfig(dslText string, cliArgs []string) Resolved

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 RetryPolicy struct {
	MaxAttempts int    `json:"max_attempts,omitempty" yaml:"max_attempts,omitempty"`
	Backoff     string `json:"backoff,omitempty" yaml:"backoff,omitempty"`
}

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 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 NewState

func NewState(values map[string]any) *State

func NewStateFromDAG

func NewStateFromDAG(dagState dag.ReadState) *State

func (*State) Get

func (s *State) Get(path string) (any, bool)

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 NewWorkflowTest(t testing.TB, reg *action.Registry, dsl string) *WorkflowFixture

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

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
Log nodes.
Log nodes.
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).
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
flow/testkit/flowtest.go
flow/testkit/flowtest.go
Package transport is the small integration boundary between Flow and concrete transport libraries.
Package transport is the small integration boundary between Flow and concrete transport libraries.

Jump to

Keyboard shortcuts

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