README
¶
pkg/core — flow runner
A small, typed engine that executes a flow (a DAG of nodes) — HTTP requests, assertions, branches, loops, polling, SSE, sub-flows — and returns a structured per-node result.
It is a greenfield rewrite of the old pkg/engine/pkg/node runner, built for
three goals:
- Simple to extend — adding a node, an assertion operator, or an output is a few lines in one file. No engine changes.
- Typed — the only place dynamic (
any) data lives isvalue.Value. Every other package is statically typed and touches runtime data through it. - No god-context — a node receives an explicit, narrow
Runtimeof the effects it may use (HTTP, Clock, Subflow, Vars). Nothing else.
Quick start
import (
"context"
"github.com/nanostack-dev/echopoint-runner/pkg/core/engine"
"github.com/nanostack-dev/echopoint-runner/pkg/core/flow"
"github.com/nanostack-dev/echopoint-runner/pkg/core/nodes" // registers the built-in node kinds
"github.com/nanostack-dev/echopoint-runner/pkg/core/value"
)
f, err := flow.Parse([]byte(`{
"name": "demo",
"nodes": [
{"id": "login", "type": "request", "method": "POST", "url": "https://api/login",
"outputs": [{"name": "token", "path": "body.access_token"}]},
{"id": "me", "type": "request", "url": "https://api/me",
"headers": {"Authorization": "Bearer {{login.token}}"},
"assertions": [{"path": "status", "op": "equals", "expected": 200}]}
],
"edges": [{"source": "login", "target": "me"}]
}`))
if err != nil { /* ... */ }
eng := engine.New(nodes.DefaultRuntime(), nil) // resolve=nil: no sub-flows
res := eng.RunFlow(context.Background(), f, value.Map{})
fmt.Println(res.Success) // true/false for the whole run
fmt.Println(res.Nodes["me"].Status) // success | failed | skipped
Registration is import-triggered. The built-in kinds register themselves in
init(). Blank-importpkg/core/nodes(as above, orimport _ ".../pkg/core/nodes") so the registry is populated before you parse a flow.
Mental model — three things
value.Value the ONLY place `any` lives. a box of decoded JSON + typed getters.
the node seam a node = a typed Cfg (embeds node.Base) + one Run function.
the engine orchestration ONLY: schedule, assert/output, recurse. no per-kind logic.
Every other package serves one of these.
Package layout
| Package | Role |
|---|---|
value |
value.Value — the dynamic-data boundary. JSONPath Get, typed getters. |
assert |
Declared assertions: a Spec (path, op, expected) evaluated over a value. |
output |
Declared outputs: bind a name to a JSONPath into a value. |
dynamicvars |
{{$gen:args}} dynamic-variable generators. |
flow |
The parsed graph: pure data. Imports neither node nor engine. |
tmpl |
Resolves {{ref}} / {{{ref}}} / {{$dyn}} in a node's raw JSON. |
node |
The authoring seam: Base, Runtime, Result, Register, capabilities. |
nodes |
The nine built-in node kinds (one file each) + DefaultRuntime. |
engine |
The scheduler + assert/output post-step + sub-flow recursion + middleware/observer. |
result |
The outcome types: FlowResult, NodeResult, statuses, skip reasons. |
Dependencies point strictly upward: value ← assert/output/tmpl ← node ←
nodes / engine. No cycles.
Execution lifecycle
What RunFlow does, in order:
JSON ─Parse→ Flow ─validate→ schedule ─loop→ [ classify → step → runNode ] → FlowResult
│
template → decode → Run → assert/output → middleware
│
(composite Run → recurse into the engine)
- Parse (
flow.Parse) — text → a graph ofNode{ID, Kind, Raw}+Edges. Each node's full object is kept as opaqueRawbytes;flownever learns what a kind means. - Validate (
engine.validateFlow) — edges reference real nodes; every Router node's targets have an edge; every FlowReferencer's child flow exists and the reference graph is acyclic. All generic — no per-kind knowledge. - Schedule (
engine.schedule) — Kahn topological order. Roots (indegree 0) seed a ready queue; finishing a node releases its successors. An outputstore(map[nodeID]→outputs, with""holding flow inputs) grows as nodes complete. - classify (
engine.classify) — per node, decide run vs skip from predecessor state: a live succeeded edge → run; a failed/skipped/routed-away predecessor → skip with a reason;run_when: alwaysnodes run for cleanup even after a failure. - step — run or skip, record a
NodeResult, emit an event, release successors. A failure is recorded, not thrown — dependents skip, the rest of the flow continues. - runNode — the universal per-node pipeline:
- template the raw JSON against the input view (
tmpl.Resolve), - decode the resolved bytes into typed config (
node.Decode), - run the node (
Bound.Run), - for a provider node, assert + extract outputs uniformly. Steps 3–4 are wrapped as one unit by any middleware (retry re-asserts too).
- template the raw JSON against the input view (
- recurse —
module/loop/pollcall back into the engine through the injectedSubflowdependency; the engine is its ownSubflowRunner.
Node catalog
All nine kinds. type is the JSON discriminator.
type |
Purpose | Key fields | Outputs |
|---|---|---|---|
request |
HTTP request; exposes the response for assertions/outputs | method, url, headers, body, timeout_ms |
{status, headers, body} |
set_variable |
Compute named values from templates | variables |
the computed map |
assert |
Validate already-produced data by fully-qualified path | assertions (on Base) |
— |
branch |
Route to one successor by the first matching case | cases[]{when, target}, default |
{matched, matched_index} |
loop |
Foreach an inline body flow; aggregate iteration outputs | items, body, item_var, index_var, max_iterations, continue_on_error |
{results, count} |
poll |
Re-run a body until the exit assertions pass (or budget) | body, assertions, max_attempts, interval_ms, timeout_ms |
{attempts, result} |
sse |
Consume text/event-stream; assert per event until a stop |
url, headers, max_events, completion_event, stop_on_assertion_failure, timeout_ms |
{events, count, last, stop_reason} |
module |
Run a named child flow once with inputs | body_flow_id, inputs |
the child's outputs |
delay |
Sleep (cancellation-aware) | duration_ms |
{delayed_ms} |
Every node also accepts the common Base fields: id, display_name,
run_when (on_success default, or always), assertions, outputs.
Result modes
A node's Run returns a node.Result shaped one of four ways — this is how the
engine knows what post-processing to apply:
- provider (
request,set_variable,loop,assert): setsProvided:trueand anAssertvalue → the framework runs the node's declared assertions/outputs. - self-evaluating (
poll,sse): evaluates its own assertions internally and returns them onAssertions→ the framework records but does not re-run them. - routing (
branch): returnsRouted(the chosen successors) → the engine marks the not-taken edges dead. - plain (
delay,module): justOutputs.
Assertions, outputs, templating
All three address data with the same path syntax (value.Value.Get):
- A bare dotted path —
body.access_token,headers.content-type,create-user.id(hyphens fine) — resolves member-by-member. - A path starting with
$is full RFC-9535 JSONPath —$.items[*].id,$.data[?@.active].
Assertion operators (assert.Op):
| Group | Operators |
|---|---|
| Equality | equals, not_equals |
| String | contains, not_contains, starts_with, ends_with, regex |
| Presence | empty, not_empty, exists |
| Numeric | gt, lt, gte, lte, between (expected: [min, max]) |
Equality/substring operators are deliberately lenient — they compare
stringified forms, so 200 == "200". Numeric operators compare real numbers
(fractional values are not truncated).
Templating (tmpl), resolved before a node is decoded, so a node never sees
a template:
| Form | Meaning |
|---|---|
{{ref}} |
inline string interpolation |
{{{ref}}} |
whole-value substitution (preserves object/number/bool type) |
{{$name:a:b}} |
a dynamic-variable generator (see dynamicvars) |
ref is the same path syntax: login.token, $.items[0].id, or a bare flow-input
name. Unresolved refs are left verbatim so a typo is visible, not silently empty.
Result model
RunFlow returns a *result.FlowResult — node failures are recorded here, never
returned as a Go error.
type FlowResult struct {
Success bool // false if any on_success node failed, or validation failed
Nodes map[string]*NodeResult // per-node outcome
Outputs value.Map // every node's outputs, nested under its id
Error string // first main-phase failure (or validation/cancel)
Code string
}
type NodeResult struct {
ID, Kind string
Status Status // success | failed | skipped
Outputs value.Map
Assertions assert.Results // per-assertion outcomes
Error, Code, SkipReason string
}
Skip reasons (a wire contract — exact strings must not drift):
dependency_failed, dependency_skipped, routed_away_by_branch,
aborted_after_failure, missing_inputs.
Error codes — a stable taxonomy. User-caused failures carry a code; a genuine
runner fault is RUNNER_ERROR.
| Code | When |
|---|---|
FLOW_VALIDATION_FAILED |
empty flow, bad edge, cycle, unreachable nodes |
INVALID_NODE_CONFIG |
template or decode error on a node |
ASSERTION_FAILED |
a declared assertion failed |
UNKNOWN_REFERENCE |
assert/branch referenced an unexecuted/unknown node |
REQUEST_FAILED |
HTTP build/transport/read error |
LOOP_FAILED |
items not a list, body parse error, or an iteration failed |
POLL_FAILED / POLL_BODY_FAILED / POLL_TIMEOUT / POLL_CONDITION_NOT_MET |
poll setup / body / deadline / budget exhausted |
SSE_FAILED |
connect, non-2xx, read, or assertion failure on the stream |
MODULE_FLOW_NOT_FOUND / MODULE_CYCLE_DETECTED / MODULE_FAILED |
sub-flow resolution / recursion / child failure |
SUBFLOW_FAILED |
inline body (loop/poll) failed |
CANCELLED |
the context was cancelled mid-run |
RUNNER_ERROR |
an unexpected non-user fault |
On a sub-flow failure the child's real code and message propagate up, so a loop/poll/module surfaces which inner node failed, not a generic wrapper.
Extending
Add a node kind
One file. No engine change. Write a config struct that embeds node.Base, a Run
function, and register it:
package nodes
type EchoCfg struct {
node.Base
Message string `json:"message"`
}
func runEcho(_ context.Context, cfg EchoCfg, _ value.Value, _ node.Runtime) (node.Result, error) {
return node.Result{Outputs: value.Map{"echo": value.Of(cfg.Message)}}, nil
}
func init() { node.Register(spi.KindEcho, runEcho) } // add KindEcho to pkg/spi/kind.go
- Embedding
node.Baseis what lets you register — the seam is sealed (onlyBase-embedders satisfy the constraint). - Need effects? Take them off
rt node.Runtime(rt.HTTP,rt.Clock,rt.Subflow,rt.Vars). Need nothing? Ignore it. - Want the framework to run assertions/outputs against your result? Set
Provided: trueand anAssertvalue. - Want validation to see references or route targets? Implement
node.FlowReferencerornode.Router— the engine picks them up generically.
Add an assertion operator
Add the constant and a case in assert.compare (assert/assert.go). Nothing
else.
Add an output
Outputs are declarative already — {"name": ..., "path": ...} on any provider
node. No code needed.
Runtime, middleware, observer
// The explicit effect set. Build one with only the fields your flow needs.
type Runtime struct {
HTTP HTTPDoer // request/sse
Clock Clock // delay/poll (WallClock in prod; fakeable in tests)
Subflow SubflowRunner // injected by engine.New — do not set
Vars DynamicResolver // {{$dyn}} generators
}
eng := engine.New(nodes.DefaultRuntime(), resolveChildFlow,
engine.WithMiddleware(engine.Retry(3), engine.Timeout(10*time.Second)),
engine.WithObserver(func(ev engine.Event) { /* progress streaming */ }),
)
- Middleware wraps each node's run-and-assert unit, so
Retryre-runs the assertions too;Retrywaits a small ctx-respecting backoff between attempts. - Observer receives
NodeStarted/NodeCompleted/NodeFailedand the flow-level events, for the top-level flow only — a node running a sub-flow emits as a single node event; its inner nodes are silent (keeps the wire flat).
Testing
gofmt -w pkg/core/
go test -race -count=1 ./pkg/core/...
golangci-lint run --path-mode=abs --timeout 5m ./pkg/core/...
Each leaf package has focused unit tests; engine/engine_test.go drives whole
flows against httptest servers with a fakeClock, covering every node kind, the
skip cascade, routing, sub-flow recursion, and the middleware.
Benchmarks
engine/bench_test.go stresses the engine with generated graphs. To keep the
measurement on engine overhead rather than the network, HTTP/SSE are served by
an instant in-memory node.HTTPDoer (canned responses, zero latency) and time by
a fakeClock. A separate BenchmarkRealisticHTTP drives a real httptest server
(the "wiremock" end-to-end sanity check).
go test ./pkg/core/engine/ -run '^$' -bench . -benchmem
Graph shapes cover every node kind: WideDiamond (root → N parallel requests →
sink asserting over all N), DeepChain (N templated requests in series), Loop
(foreach over N items), SSE (N-event stream), VarChain (N set_variables —
pure engine, no HTTP), BranchFanout (route past N cascade-skipped successors),
Modules (N parallel sub-flows), Poll, DelayChain, and Complex (a
kitchen-sink graph combining request + branch + loop + module + assert).
Eleven optimizations came out of profiling these:
- Incremental input view.
runNodeused to rebuild the whole output store into a fresh map on every node (inputView), re-boxing each node's outputs each time — O(n²) allocation, ~74% of all allocs on a wide graph. The scheduler now maintains the denormalized view incrementally (publish each node's outputs once on completion; box in O(1) per step). Quadratic → linear. - Capability-indexed validation.
validateTargets/walkRefsdecoded every node just to check whether it routes or references a flow. Each kind's capabilities are now probed once atRegister(node.Routes/node.References), so validation skips decoding kinds that can never route/reference. - Template fast-path + parsed
run_when.tmpl.Resolveunmarshalled + walked + re-marshalled every node even with no{{tokens; it now returns raw untouched when there are none.run_whenis lifted intoflow.Nodeat parse time instead of re-unmarshalled per node. Together: static nodes skip two JSON round-trips (e.g. DelayChain/BranchFanout ~−57% allocs). - JSONPath parse cache.
value.Getre-parsed the path expression on every call; the same assertion/template paths repeat across iterations and runs, so async.Mapmemoizes the parsed*jsonpath.Path. Biggest win on eval-heavy flows — Loop 1000 items ~−48% allocs. - Lean scheduler state. Per run the scheduler allocated nine maps;
done+failedare merged into onestatemap,deadis lazy (branch-free flows never allocate it), and the topology maps (indeg/succ/preds) are skipped entirely for an edgeless body. Sub-flow-heavy flows run many schedulers, so this compounds (Loop, Complex ~−6%). - Per-string template gate.
resolveStringran two regexes on every string in a templated node, including static ones. Astrings.Contains(s, "{{")gate skips both for strings with no token (DeepChain ~−9% allocs). - No empty allocations.
assert.Run/output.Extractreturned an empty slice/map for nodes that declare none, andvalidateTargetsbuilt an edge index for flows with no routing node — all now short-circuit to nil. - Dotted-path fast walker. A bare dotted path ("node.key" — the
overwhelmingly common case) used to be bracket-quoted, parsed (cached), and
evaluated through the JSONPath engine, which allocates per
Select.value.Getnow walks bare dotted paths member-by-member with zero allocations; only "$"-prefixed full JSONPath goes through the library. (WideDiamond/DeepChain time ~−20%.) - Byte-level template rewrite.
tmpl.Resolveunmarshalled the whole node config toany, walked it, and re-marshalled — thennode.Decodeunmarshalled the result again (three JSON passes for a templated node). Templates can only occur inside JSON string values, so Resolve now scans the raw bytes, rewrites only string literals that contain a token (object keys and token-free strings are copied verbatim), and leaves one JSON pass: the typed decode. (VarChain/DeepChain ~−30% time, ~−40% allocs.) - One boxing per node. Provider nodes boxed their outputs for the
assertion post-step (
Assert: out.Value()) even when the node declares no assertions or outputs; the engine boxed the same map again for the view and a third time incollect. Now: a provider leavesAssertzero and the engine boxes outputs once — and only when assertions/outputs are declared — andcollectreuses the view's boxing. - Decode-once values.
set_variable/moduleheldmap[string]json.RawMessageand re-parsed every entry withvalue.JSONat run time, andassert.Spec.Expectedstayed raw JSON and was re-parsed on every evaluation (per poll attempt, per SSE event, per loop iteration).value.ValuegainedUnmarshalJSON, so these fields decode straight intovalue.Value/value.Mapin one pass. (SSE −27% time, −34% allocs.)
(A twelfth candidate — caching each flow's topology at parse — was prototyped
and reverted: it added derived state + a dual code path to the deliberately
pure-data flow package for only a ~7% loop/module gain, not worth the surface.)
Cumulative vs the pre-optimization baseline, WideDiamond n=256: time −87%,
memory −90%, allocs −85%, and −52% time / −58% allocs on the pure-engine
VarChain n=256. Beyond this the remaining cost is inherent per-node work — the
single typed json.Unmarshal of each node's config, evaluating "$" JSONPath, and
boxing each node's outputs into the pure-any view — not waste.
Design principles (the invariants worth keeping)
anylives only invalue. Every other package is statically typed.- The engine has no per-node-type logic. Every kind dispatches through the same registry; branch/module specifics are surfaced via capability interfaces.
- Nodes are tiny (30–90 lines). Scheduling, templating, asserting, skipping and
recursing live once in the engine; a node just returns a
Result. - Explicit dependencies, no god-context.
context.Contextis for cancellation; effects come fromRuntime. - Failure-continue. One node failing records an outcome and skips dependents; the caller always sees the whole run.
Directories
¶
| Path | Synopsis |
|---|---|
|
Package assert evaluates declared assertions against a value.
|
Package assert evaluates declared assertions against a value. |
|
Package dynamicvars resolves {{$name:args}} template variables to generated fake data.
|
Package dynamicvars resolves {{$name:args}} template variables to generated fake data. |
|
Package engine is orchestration only.
|
Package engine is orchestration only. |
|
Package flow is the parsed graph: pure data, no behavior.
|
Package flow is the parsed graph: pure data, no behavior. |
|
Package node is the node-authoring seam.
|
Package node is the node-authoring seam. |
|
Package output extracts declared named outputs from a value.
|
Package output extracts declared named outputs from a value. |
|
Package result is the outcome of a flow run: a per-node record plus the flow-level verdict.
|
Package result is the outcome of a flow run: a per-node record plus the flow-level verdict. |
|
Package tmpl resolves {{ref}} / {{{ref}}} template tokens in a raw node definition against the node's input view (flow inputs + upstream outputs) and optional dynamic-variable generators, before the node is decoded.
|
Package tmpl resolves {{ref}} / {{{ref}}} template tokens in a raw node definition against the node's input view (flow inputs + upstream outputs) and optional dynamic-variable generators, before the node is decoded. |
|
Package value boxes decoded-JSON data behind a typed accessor API.
|
Package value boxes decoded-JSON data behind a typed accessor API. |