nmbrs-runtime 0.4.0

Workload execution runtime for nmbrs
Documentation

nmbrs-runtime

The workload execution runtime for nmbrs. It takes a parsed workload, compiles its Polydat bindings into a tree of scope kernels, walks the scenario tree, runs phases as activities of concurrent fibers, dispatches ops through adapters and wrapper stacks, and records metrics into a session. It also defines the adapter API (DriverAdapter / OpDispenser) that protocol drivers implement. That API is the main reason to depend on this crate directly.

Where it sits in nmbrs

End users normally install the nmbrs CLI (cargo install nmbrs) and run workloads with nmbrs run. This crate is for people who write adapters, or who embed the runner in their own binary.

Writing an adapter

Adapters live in nmbrs_runtime::adapter. An adapter has two phases:

  • Init time. DriverAdapter::map_op is called once per op template, before any cycle runs. It does the expensive work: validating fields, preparing statements, building binders. It returns a boxed OpDispenser.
  • Cycle time. OpDispenser::execute(cycle, ctx) is called for every cycle. It resolves that cycle's values and performs the operation.

DriverAdapter

Implement this trait on your adapter type. It is constructed once per activity and shared across fibers through an Arc.

Method Required Purpose
name(&self) -> &str yes The adapter name, e.g. "http".
map_op(&self, template: &ParsedOp, parent: Arc<dyn Kernel>) -> MapOpFuture yes Build the dispenser for one op template. parent is the phase's Polydat scope kernel.
known_op_fields() no Return Some(&[...]) to declare your op-field vocabulary. The core then rejects templates with unknown fields. None (the default) is permissive.
known_op_params() no Extra keys allowed under an op's params.
default_status_metrics() no StatusMetrics shown on the status line.
display_preference() no DisplayPreference::Off if the adapter writes to the raw terminal. The default is Auto.
declare_controls(parent) no Declare adapter-level dynamic controls on a subcomponent of the activity's nmbrs_metrics component.
shutdown() no Async teardown when a shared adapter is released.
accessor_payload() no A type-erased handle that kernel nodes can look up through the resource scope.

If map_op returns Err, construction stops before any cycle runs. If a dispenser binds fields through a typed API, and not by plain text substitution, it should build a Binder and verify it against parent (adapter::verify_binders) inside map_op. Kernel, Binder, BinderSlot, PortType and ExecCtx are re-exported from nmbrs_runtime::adapter, so an adapter crate doesn't need a direct polydat dependency.

OpDispenser

Method Required Purpose
execute(&self, cycle, ctx: &ExecCtx) -> Pin<Box<dyn Future<Output = Result<OpResult, ExecutionError>> + Send>> yes Run one op.
canonical_kernel() no Return Some(&kernel) when the dispenser keeps the Polydat kernel it got from map_op. The executor builds per-fiber kernels from it.
describe() / describe_resolved(wires) no One-line views of the op, as a template and as sent. Used in error diagnostics.
adapter_metrics() no Extra (family, Labels, MetricValue) samples to include in metrics snapshots.
status_counters() no Cumulative (name, count) pairs for the status line. A leading _ marks a name as internal.
rows_per_op() no The cursor stride for batch ops. The default is 1.
inner_dispenser() no Wrappers must return their inner dispenser. Leaves keep the default None.

At cycle time, ctx.wires (a wires::WireSource) is how a dispenser reads bound values. It resolves names against the fiber's kernel for this dispenser. wires::resolve_op_fields_via_wires renders a list of op fields:

  • A field that is exactly {name} keeps its typed value.
  • References embedded in text are substituted.
  • Bare strings stay literal.

The result is a ResolvedFields, which has names, typed values, get_str, get_value and strings().

Results and errors:

  • Success is an OpResult. Its body is an Option<Box<dyn ResultBody>>. TextBody and JsonBody are provided. Implement ResultBody for native result types; to_json, element_count and byte_count are used for captures and traversal metrics.
  • Failure is an ExecutionError:
    • ExecutionError::Op(AdapterError) is a failure of this op.
    • ExecutionError::Adapter(AdapterError) means the connection or session is degraded.
    • AdapterError has three fields. error_name is the name that the workload's errors: rules match against. message is shown to the user. retryable: true marks an Op error that the tries: wrapper may re-run.

Registration

Adapters register at link time with inventory, by submitting an AdapterRegistration:

  • names are the names users select with adapter=<name>.
  • known_params are adapter params, used for CLI validation.
  • display_preference decides TUI compatibility from the params.
  • supported_controls is a list of control_catalog::ControlDesc descriptors.
  • create is an async factory from params to an Arc<dyn DriverAdapter>.

There are two optional registrations:

  • DriverImpl lets one adapter have several driver implementations. The user picks one with a selector param, such as cqldriver=.
  • SharedDriverRegistration lets phases with the same ResourceKey share one adapter instance through the resource pool.

A minimal adapter that renders its op fields as text:

// Cargo.toml: nmbrs-runtime = "0.3", nmbrs-workload = "0.3",
//             inventory = "0.3", serde_json = "1"
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;

use nmbrs_runtime::adapter::{
    AdapterError, AdapterRegistration, DisplayPreference, DriverAdapter, ExecCtx,
    ExecutionError, Kernel, MapOpFuture, OpDispenser, OpResult, TextBody,
};
use nmbrs_runtime::wires::resolve_op_fields_via_wires;
use nmbrs_workload::model::ParsedOp;

struct EchoAdapter;

impl DriverAdapter for EchoAdapter {
    fn name(&self) -> &str {
        "echo"
    }

    fn map_op<'a>(&'a self, template: &'a ParsedOp, parent: Arc<dyn Kernel>) -> MapOpFuture<'a> {
        Box::pin(async move {
            // Init-time: snapshot the op fields once.
            let fields: Vec<(String, serde_json::Value)> =
                template.op.iter().map(|(k, v)| (k.clone(), v.clone())).collect();
            Ok(Box::new(EchoDispenser { kernel: parent, fields }) as Box<dyn OpDispenser>)
        })
    }
}

struct EchoDispenser {
    kernel: Arc<dyn Kernel>,
    fields: Vec<(String, serde_json::Value)>,
}

impl OpDispenser for EchoDispenser {
    fn canonical_kernel(&self) -> Option<&Arc<dyn Kernel>> {
        Some(&self.kernel)
    }

    fn execute<'a>(
        &'a self,
        _cycle: u64,
        ctx: &'a ExecCtx<'a>,
    ) -> Pin<Box<dyn Future<Output = Result<OpResult, ExecutionError>> + Send + 'a>> {
        Box::pin(async move {
            let resolved = match resolve_op_fields_via_wires(&self.fields, ctx.wires) {
                Ok(r) => r,
                Err(message) => {
                    return Err(ExecutionError::Op(AdapterError {
                        error_name: "BindError".into(),
                        message,
                        retryable: false,
                    }));
                }
            };
            Ok(OpResult {
                body: Some(Box::new(TextBody(resolved.strings().join(" ")))),
                skipped: false,
            })
        })
    }
}

inventory::submit! {
    AdapterRegistration {
        names: || &["echo"],
        known_params: || &[],
        display_preference: |_params| DisplayPreference::Auto,
        supported_controls: || &[],
        create: |_params| Box::pin(async move {
            Ok(Arc::new(EchoAdapter) as Arc<dyn DriverAdapter>)
        }),
    }
}

Reference implementations in the repository:

  • stdout is the smallest real adapter. It also shows SharedDriverRegistration.
  • testkit injects errors deterministically.
  • http and cql are full adapters. cql uses DriverImpl for its scylla and cassandra-cpp drivers.

Running with your adapter

The published nmbrs binary only contains the adapters it was built with. To use your own adapter, build a binary that links your adapter crate and calls the runner. Link the crate explicitly (for example with extern crate), so that its inventory::submit! registration is kept:

extern crate my_adapter; // force-link for inventory registration

#[tokio::main]
async fn main() {
    let args: Vec<String> = std::env::args().skip(1).collect();
    if let Err(e) = nmbrs_runtime::runner::run(&args).await {
        eprintln!("error: {e}");
        std::process::exit(1);
    }
}
my-bench run workload=my_workload.yaml adapter=echo cycles=10

runner::run accepts run or a bare key=value / workload-file argument list. Other nmbrs subcommands (check, report, …) are implemented in the nmbrs crate, not here.

Execution model

  • Runner (runner). run(args) and run_with_observer(args, observer) run the whole pipeline: parameter handling, workload resolution and parsing, adapter creation, scenario execution, metrics and the session. run_executions runs several executions in one shared session. concurrent::run_workload_headless runs one workload and returns an ExecutionOutcome.
  • Scope tree (scope_tree, scope_kernel, scope, bindings). Workload, scenario and phase bindings compile into a ScopeTree of ScopeKernels. Children are bound under their parent's kernel, and op templates' kernels are bound under the phase kernel they get in map_op.
  • Scenario executor. It is internal. It walks the ScenarioNode tree at run time and evaluates for_each / do_while / do_until as it goes. scene_tree::SceneTree is the view of that walk that renderers use, with iterations unrolled into per-iteration phase nodes.
  • Activities and fibers (activity). Each phase runs as an Activity, configured by ActivityConfig:
    • concurrency fibers (tokio tasks) execute stanzas.
    • opseq maps cycles to ops by ratio (bucket, interval or concat sequencing).
    • An optional activity-level nmbrs_rate::RateLimiter gates all fibers.
    • Stop conditions, throttle: and tries: settings apply per activity.
  • Wrappers (wrappers, wrapper_registry, wrapper_resolver). Op fields such as tries:, if:, delay:, poll:, rate:, metrics: and errors: select wrapper dispensers. These compose around the adapter's dispenser in a fixed innermost-to-outermost order (see wrapper_resolver::DEFAULT_ORDER). Wrappers implement adapter::WrappingDispenser and expose inner_dispenser().
  • Error routing. The outermost op wrapper applies the op's nmbrs_errorhandler::ErrorRouter policy to each terminal failure.
  • Observers (observer). RunObserver receives phase and op lifecycle events. Implementations:
    • StderrObserver, the plain-text default
    • concurrent::HeadlessObserver, which collects an outcome with no display
    • TuiObserver, provided by nmbrs-tui
  • Metrics. Each activity registers its instruments (ActivityMetrics) on a component in the nmbrs_metrics component tree. Adapter samples from OpDispenser::adapter_metrics are added to the same snapshots. Dynamic controls such as concurrency and rate are declared on the activity component. Each session writes its metrics to a SQLite metrics.db in its session directory (session).

Cargo features

Feature Default Enables
flamegraph no The profiler=flamegraph mode: in-process CPU sampling with pprof. The SVG is rendered by the external inferno-flamegraph tool when it is on PATH. Without this feature, profiler=flamegraph logs a warning and does nothing.

The nmbrs CLI forwards its own flamegraph feature to this one.

Links

License

Apache-2.0