Python¶
Built with pybind11. Install with pip.
python -m pip install flowedge-0.1.2-*.whl # CPU wheel from the v0.1.2 GitHub Release
# or, from a clone with a C++23 toolchain:
python -m pip install .
Wheels are CPU. CUDA is a source build. PyPI is not published. Python quickstart.
import numpy as np, flowedge
e = flowedge.Engine("models/mamba_flow.safetensors")
print(e.action_dim, e.condition_dim, e.d_model, e.thread_count, e.cuda_resident)
print(e.model_metadata) # architecture, precision, dimensions, digest, snapshot size
tokens = np.array([1, 2, 3, 4], dtype=np.int32)
hidden = e.run(tokens)
noise = np.zeros(e.action_dim, dtype=np.float32)
action = e.sample(prefix=tokens,
noise=noise,
steps=10, method="euler")
The convenience calls accept array-like inputs, normalize their dtype/layout when necessary, and
allocate returned NumPy arrays. A real-time loop can instead own every output and use the zero-copy
forms; those inputs and outputs must expose C-contiguous int32 or float32 buffers:
hidden = np.empty((tokens.size, e.d_model), dtype=np.float32)
e.run_into(tokens, hidden)
current = np.empty(e.d_model, dtype=np.float32)
e.step_into(7, current)
action = np.empty(e.action_dim, dtype=np.float32)
e.sample_into(tokens, noise, action, steps=10, method="euler")
All inference calls release the GIL while C++ runs.
Pass threads=0..8 to override the automatic worker pool for a specific engine. If omitted,
FLOWEDGE_THREADS is honored and then the automatic default is used.
e.cuda_resident is runtime residency, not flowedge.cuda (compile-time).
CPU wheels always report false. A CUDA source build reports false when attach
failed and the engine is on CPU kernels. FLOWEDGE_CUDA_REQUIRED=1 fails load
in that case.
single_threaded = flowedge.Engine("models/mamba_flow.safetensors", threads=0)
External encoders and cooperative solving¶
A head-only checkpoint accepts a condition vector produced by PyTorch, TensorRT, ONNX Runtime, a shared-memory camera process, or another model server. Output is caller-owned:
condition = encoder(observation).astype(np.float32, copy=False)
noise = np.zeros(e.action_dim, dtype=np.float32)
action = np.empty(e.action_dim, dtype=np.float32)
e.sample_condition(condition, noise, action, 8, "heun")
For a generic Transformer checkpoint, the experimental boundary accepts
caller-owned embedding rows. A batch-one prefix mask may be supplied for padded
VLA sequences; it must contain
leading 1 values followed by 0 padding:
transformer = flowedge.Engine("models/tiny-gpt2.flowedge.safetensors")
embeddings = external_vlm(observation).astype(np.float32, copy=False)
mask = np.array([1, 1, 1, 0], dtype=np.uint8)
hidden = transformer.run_embeddings(embeddings, attention_mask=mask)
This is an embedding boundary, not a native SmolVLA encoder.
The pinned SmolVLA checkpoint exposes the native action expert when the caller supplies a captured VLM K/V cache. Its three allocation-free primitives are suffix embedding, expert execution, and action projection:
smolvla = flowedge.Engine("models/smolvla_base/model.safetensors")
noisy_actions = np.zeros((50, 32), dtype=np.float32)
suffix = smolvla.smolvla_embed_suffix(noisy_actions, timestep=1.0)
assert suffix.shape == (50, 720)
# Cache layout: [16 expert layers, prefix length, 320 VLM K/V values].
# Keys have already received the VLM RoPE transform. The uint8 mask has one
# validity byte per prefix slot; source SmolVLA may have zero-padded language
# tokens before a valid state token.
hidden = smolvla.smolvla_run_expert(suffix, prefix_keys, prefix_values, prefix_mask)
velocity_padded = smolvla.smolvla_project_actions(hidden)
# Prefer this equivalent single call in the control loop.
velocity_padded = smolvla.smolvla_denoise(
noisy_actions, 1.0, prefix_keys, prefix_values, prefix_mask
)
# Source SmolVLA's deterministic 10-step Euler loop, seeded by the caller.
actions_padded = smolvla.smolvla_sample(
initial_noise, prefix_keys, prefix_values, prefix_mask, steps=10
)
This does not run LeRobot observation preprocessing, image/language encoding, or build the VLM cache. A real captured-cache trajectory replay artifact is still required before making a native SmolVLA inference or control-quality claim.
The same solve can be split across scheduler quanta without changing its result:
e.flow_begin(condition, noise, 8, "heun", generation=42,
timestamp_ns=observation_time, deadline_ns=control_deadline)
remaining = e.flow_advance(action, 2)
while remaining:
remaining = e.flow_advance(action, 1)
Only publish action as final when remaining == 0. One engine owns one resumable solve.
e.flow_metadata reports the generation, original timestamp/deadline, status, model digest, and
remaining NFE. A scheduler can call e.cancel_before(43) concurrently; the older solve stops at
the next complete solver-step boundary and raises RuntimeError from flow_advance.
Streaming state¶
step(token) advances the Mamba recurrence, reset() starts a fresh stream, and
decode_state() / restore_decode_state(bytes) create deterministic branches. Snapshots are
versioned and checksummed; restore rejects a different model digest, architecture, precision,
dimensions, truncated payload, or corruption before changing state.
Source: python/flowedge_ext.cc.
Diffusion Policy¶
e = flowedge.Engine("models/diffusion_pusht.flowedge.safetensors")
condition = observation_encoder(history).astype(np.float32, copy=False)
noise = np.random.default_rng(7).standard_normal(
(e.action_horizon, e.action_dim), dtype=np.float32)
actions = e.sample_diffusion(condition, noise, steps=10, scheduler="ddim")
actions is [action_horizon, action_dim] and is already in dataset action
units. sample_diffusion_into writes into a caller-owned C-contiguous float32
buffer. Seeded scheduler="ddpm" is deterministic; the DDIM result depends only
on the condition and supplied initial noise. diffusion_denoise runs one
normalized U-Net pass for verification or profiling.
The LeRobot companion adapter exposes the same allocation-conscious pattern through
predict_action_chunk_into and select_action_into; reuse those buffers in a control loop.
e.diffusion_metadata exposes horizon, action_steps, observation_steps,
train_timesteps, and the clipping policy. LeRobot normally executes
actions[observation_steps - 1: observation_steps - 1 + action_steps] from the
returned horizon; the controller owns that slicing decision.