runtime

package
v0.20.0 Latest Latest
Warning

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

Go to latest
Published: Oct 8, 2026 License: Apache-2.0 Imports: 19 Imported by: 0

Documentation

Overview

Package runtime ships the primitive actions every .nflow pipeline starts from: runtime.const, runtime.noop, runtime.debug, runtime.fail, runtime.pick, runtime.wrap, runtime.with, runtime.env, runtime.uuid, runtime.call, runtime.dispatch_by_prefix, json.clean, runtime.sleep. This is a native bundle — always mounted by native.Bundles().

Typical use:

{ user_id: 42 } -> runtime.pick @{ field: "user_id" } -> runtime.wrap @{ key: "data" }
runtime.fail @{ kind: "Timeout", message: "upstream" } || runtime.const @{ value: "fallback" }

Index

Constants

View Source
const ID = "runtime"

Variables

View Source
var Call = action.New("runtime.call", func(ctx context.Context, in any) (any, error) {
	m, ok := in.(map[string]any)
	if !ok {
		return nil, xerr.BadRequest("call: input must be an object carrying 'name' and optional 'payload'")
	}

	name := strings.TrimSpace(readStringArg(m, "name"))
	if name == "" {
		return nil, xerr.BadRequest(`call: 'name' is required (use: { name: "log.info", payload: { ... } } -> runtime.call)`)
	}

	resolver := contracts.ActionResolverFromContext(ctx)
	if resolver == nil {
		return nil, xerr.Internal("call: no action resolver in execution context")
	}

	target, ok := resolver.Action(name)
	if !ok {
		return nil, xerr.NotFound("call: target action " + name + " not found in registry")
	}

	payload, hasPayload := m["payload"]
	if !hasPayload || payload == nil {
		filtered := make(map[string]any, len(m))
		for k, v := range m {
			if k != "name" {
				filtered[k] = v
			}
		}
		payload = filtered
	}
	return action.InvokeAny(ctx, target, payload)
}).Description("Resolve and invoke another action by name at runtime").
	Tag("base", "dynamic").
	Build()
View Source
var Const = action.New("runtime.const", func(_ context.Context, in any) (any, error) {
	if s, ok := in.(string); ok {
		return coerceLiteral(s), nil
	}

	m, ok := in.(map[string]any)
	if !ok {
		return nil, xerr.BadRequest("const: expected object with 'value' key or raw scalar, got " + typeName(in))
	}

	raw, ok := m["value"]
	if !ok {
		raw, ok = m["val"]
	}
	if !ok {
		return nil, xerr.BadRequest("const: missing required field 'value'")
	}
	return coerceLiteral(raw), nil
}).Description("Return a fixed literal; numbers, bools, and JSON auto-parse").
	Tag("base", "literal").
	Build()

Const returns a fixed literal. Numbers, bools, JSON objects, JSON arrays, and JSON-quoted strings are auto-coerced from the string representation so the DSL stays readable: runtime.const @{ value: 42 } produces int64(42), not "42".

View Source
var Debug = action.New("runtime.debug", func(_ context.Context, in any) (any, error) {
	label := readStringArg(in, "label")
	suffix := ""
	if label != "" {
		suffix = " " + label
	}

	data, err := json.MarshalIndent(in, "", "  ")
	if err != nil {
		fmt.Fprintf(os.Stderr, "[debug%s] <unprintable %T: %v>\n", suffix, in, err)
		return in, nil
	}
	fmt.Fprintf(os.Stderr, "[debug%s] %s\n", suffix, string(data))
	return in, nil
}).Description("Print input as JSON to stderr, pass through unchanged").
	Tag("base", "debug").
	Build()

Debug prints the input as JSON to stderr and returns it unchanged. The optional @{ label: "..." } arg prefixes the printed line.

View Source
var DispatchByPrefix = action.New("runtime.dispatch_by_prefix", func(ctx context.Context, in any) (any, error) {
	m, ok := in.(map[string]any)
	if !ok {
		return nil, xerr.BadRequest("dispatch_by_prefix: input must be an object")
	}

	prefix := strings.TrimSpace(readStringArg(m, "prefix"))
	key := strings.TrimSpace(readStringArg(m, "key"))
	if prefix == "" || key == "" {
		return nil, xerr.BadRequest("dispatch_by_prefix: prefix and key are required")
	}

	name := prefix + "." + key
	resolver := contracts.ActionResolverFromContext(ctx)
	if resolver == nil {
		return nil, xerr.Internal("dispatch_by_prefix: no action resolver in execution context")
	}

	target, ok := resolver.Action(name)
	if !ok {
		return nil, xerr.NotFound(fmt.Sprintf("dispatch_by_prefix: action %q is not in the registry", name))
	}

	payload := m["payload"]
	if payload == nil {
		payload = in
	}
	return action.InvokeAny(ctx, target, payload)
}).Description("Resolve `prefix.key` against the registry and invoke it").
	Tag("base", "dynamic").
	Build()
View Source
var Env = action.New("runtime.env", func(_ context.Context, in any) (any, error) {
	name := ""
	required := false

	switch v := in.(type) {
	case string:
		name = strings.TrimSpace(v)
	case map[string]any:
		name = strings.TrimSpace(readStringArg(v, "name"))
		required = isTruthyArg(v, "required")
	}

	if name == "" {
		return nil, xerr.BadRequest(`env: name is required (use: env @{ name: "HOME" })`)
	}

	value, ok := os.LookupEnv(name)
	if !ok {
		if required {
			return nil, xerr.NotFound("env: " + name + " is not set")
		}
		return "", nil
	}
	return value, nil
}).Description("Read an environment variable (required: true fails if unset)").
	Tag("base", "runtime").
	Build()

Env reads an environment variable. Accepts either a bare string (env "HOME") or an object with name and required fields.

View Source
var Fail = action.New("runtime.fail", func(_ context.Context, in any) (any, error) {
	m, _ := in.(map[string]any)
	msg := strings.TrimSpace(readStringArg(m, "message"))
	if msg == "" {
		msg = "fail: pipeline deliberately aborted by the fail action"
	}
	kind := strings.TrimSpace(readStringArg(m, "kind"))
	return nil, failError(kind, msg)
}).Description("Always return an error; kind configurable (default Internal)").
	Tag("base", "error").
	Build()

Fail always returns an error. The @{ kind: "..." } arg selects the xerr kind; the default is Internal. It is the idiomatic way to simulate a specific failure for fallback pipelines.

View Source
var Flatten = action.New("runtime.flatten", func(_ context.Context, in any) (any, error) {
	if in == nil {
		return in, nil
	}
	if _, ok := in.(map[string]any); ok {
		return in, nil
	}
	rv := reflect.ValueOf(in)
	for rv.Kind() == reflect.Pointer {
		if rv.IsNil() {
			return in, nil
		}
		rv = rv.Elem()
	}
	if rv.Kind() != reflect.Struct {
		return in, nil
	}
	data, err := json.Marshal(in)
	if err != nil {
		return in, nil
	}
	var m map[string]any
	if json.Unmarshal(data, &m) != nil {
		return in, nil
	}
	return m, nil
}).
	Description("Unwrap a struct into a map so its fields flow as top-level keys").
	Tag("base", "shape").
	Build()

Flatten unwraps a struct into a map so its fields become top-level keys in downstream @{} payloads.

Why this exists: a pipe carries typed values. When an action returns a struct (distribute.map returns DistributeMapRes, a custom action returns its own DTO), the next action's @{} args are merged on top of a wrapper map, not the struct's fields. A ref like `items` or `.items` in the following @{} atom does not reach the struct field.

Flatten solves it with one explicit step: marshal the struct to JSON, unmarshal it as a map, hand the map downstream. Now every field is a top-level key and @{} refs work as expected.

map[string]any passes through untouched. Slices, primitives, and nil pass through untouched. Only structs (and pointers to structs) are converted, because they are the only shape whose fields cannot be addressed by downstream @{} refs.

View Source
var JSONClean = action.New("json.clean", func(_ context.Context, in any) (any, error) {
	raw := ""
	var source map[string]any

	switch v := in.(type) {
	case string:
		raw = v
	case map[string]any:
		if s, ok := v["content"].(string); ok {
			raw = s
			source = v
		}
	}

	if source == nil {
		return cleanJSONString(raw), nil
	}

	out := make(map[string]any, len(source))
	maps.Copy(out, source)
	out["content"] = cleanJSONString(raw)
	return out, nil
}).Description("Strip markdown fences and prose around a JSON body").
	Tag("base", "json").
	Build()

JSONClean strips the boilerplate that reasoning models habitually wrap around structured output: markdown code fences, "Here is the JSON" preambles, trailing commentary, and single-line fence forms.

Input: string, or map with a `content` field carrying a string. Output: string (the cleaned JSON body). When input is a map, the

cleaned body replaces `content` and the map is passed through
otherwise unchanged.
View Source
var Noop = action.New("runtime.noop", func(_ context.Context, in any) (any, error) {
	return in, nil
}).Description("Pass input through unchanged").
	Tag("base", "identity").
	Build()

Noop returns the input unchanged. It is the identity action used in tests, as a placeholder in pipeline composition, and as an explicit terminator for a stream.

View Source
var Pick = action.New("runtime.pick", func(_ context.Context, in any) (any, error) {
	m, ok := in.(map[string]any)
	if !ok {
		return nil, xerr.BadRequest("pick: input must be an object")
	}

	field := strings.TrimSpace(readStringArg(m, "field"))
	if field != "" {
		val, found := lookupDottedPath(m, field)
		if !found {
			return nil, xerr.NotFound("pick: field '" + field + "' not found")
		}
		return val, nil
	}

	if onlyVal, exists := m["only"]; exists {
		keys := parseKeyList(onlyVal)
		if len(keys) > 0 {
			filtered := make(map[string]any, len(keys))
			for _, k := range keys {
				if val, found := lookupDottedPath(m, k); found {
					filtered[k] = val
				}
			}
			return filtered, nil
		}
	}

	if dropVal, exists := m["drop"]; exists {
		keys := parseKeyList(dropVal)
		if len(keys) > 0 {
			filtered := make(map[string]any, len(m))
			for k, v := range m {
				if k != "drop" && !slices.Contains(keys, k) {
					filtered[k] = v
				}
			}
			return filtered, nil
		}
	}

	return nil, xerr.BadRequest("pick: must specify 'field', 'only', or 'drop'")
}).Description("Extract field or apply allowlist/denylist to input keys").
	Tag("base", "shape").
	Build()

Pick extracts a single field or filters the map with an allowlist / denylist.

View Source
var Print = action.New("runtime.print", func(_ context.Context, in any) (any, error) {
	cfg := extractPrintConfig(in)
	data := stripPrintConfig(in)

	writer := printStream(cfg.stream)
	rendered := renderPrint(data, cfg, printIsTerminal(writer))
	writePrintOutput(writer, cfg.label, rendered)

	return data, nil
}).
	Description("Pretty-print the input with size limits; strips its own config keys before passing through").
	Tag("base", "debug", "print").
	Build()

Print renders the incoming value to stdout or stderr with a bound on depth, item count per container, and string length. It is a clean tap: its own config keys are stripped from the value that continues downstream, so a debug insertion does not change the pipeline shape.

View Source
var Sleep = action.New("runtime.sleep", func(ctx context.Context, in any) (any, error) {
	ms := 1000
	if m, ok := in.(map[string]any); ok {
		if val, exists := m["duration_ms"]; exists {
			switch v := val.(type) {
			case float64:
				ms = int(v)
			case int:
				ms = v
			}
		}
	}

	select {
	case <-time.After(time.Duration(ms) * time.Millisecond):
	case <-ctx.Done():
		return nil, ctx.Err()
	}

	return in, nil
}).Description("Pause execution for duration_ms (default 1000ms) and return input").
	Tag("base", "simulate").
	Build()

Sleep pauses execution for the specified milliseconds. It respects context cancellation, so if a pipeline times out or is canceled, the sleep aborts immediately instead of holding the goroutine.

View Source
var UUID = action.New("runtime.uuid", func(_ context.Context, in any) (any, error) {
	id, err := newUUIDv4()
	if err != nil {
		return nil, xerr.Internal("uuid: entropy source failed", err)
	}

	if key := strings.TrimSpace(readStringArg(in, "as")); key != "" {
		base := map[string]any{}
		if m, ok := in.(map[string]any); ok {
			for k, v := range m {
				if k == "as" {
					continue
				}
				base[k] = v
			}
		}
		base[key] = id
		return base, nil
	}
	return id, nil
}).Description("Generate a UUID v4; returns a string, or merges under @{ as: ... }").
	Tag("base", "runtime").
	Build()

UUID generates a UUID v4. Without args it returns a bare string; with @{ as: "field" } it merges the UUID into the input map under that key.

View Source
var With = action.New("runtime.with", func(_ context.Context, in any) (any, error) {
	return in, nil
}).Description("Merge @{...} args into the input map").
	Tag("base", "shape").
	Build()

With merges the @{...} args into the input map and passes the result downstream unchanged. Args are injected by the compiler before the handler runs, so the handler is a plain pass-through:

{ goal: "x", attempt: 1 } -> runtime.with @{ attempt: .attempt + 1 }
// → { goal: "x", attempt: 2 }

Called without @{...} it degenerates to noop. Semantics match C# record `with { ... }`: the input is the base, listed fields win.

View Source
var Wrap = action.New("runtime.wrap", func(_ context.Context, in any) (any, error) {
	m, ok := in.(map[string]any)
	if !ok {
		return nil, xerr.BadRequest("wrap: input must be an object carrying the `key` arg")
	}
	key := strings.TrimSpace(readStringArg(m, "key"))
	if key == "" {
		return nil, xerr.BadRequest(`wrap: key is required (use: wrap @{ key: "data" })`)
	}
	inner := make(map[string]any, len(m))
	for k, v := range m {
		if k == "key" {
			continue
		}
		inner[k] = v
	}
	return map[string]any{key: inner}, nil
}).Description("Wrap the input object under a named key").
	Tag("base", "shape").
	Build()

Wrap wraps the input object under a named key. The key is required: wrap @{ key: "data" }.

Functions

func Bundle

func Bundle(_ map[string]string) core.Bundle

Types

This section is empty.

Jump to

Keyboard shortcuts

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