Skip to content

Data

Stream a LeRobot dataset without a full download, sync one to a bucket, and label an episode.

Datasets are LeRobot v3 datasets on disk or on the Hub. These functions stream them without a full download, sync them to a bucket, and grade episodes.

Streaming

Streaming read-back for LeRobotDataset - read frames directly from the Hub.

Primary use: in-process eval / replay / notebooks / agent loops. Streamed training does not need this module: python -m lerobot.scripts.lerobot_train --dataset.streaming=true already uses StreamingLeRobotDataset via lerobot.datasets.factory.make_dataset.

lerobot is never imported at module top-level (numpy/pandas ABI safety on Jetson; see the :mod:strands_robots.dataset_recorder header). The [lerobot] extra floors lerobot at BUCKET_STREAMING_MIN_LEROBOT, whose constructor accepts every keyword :meth:StreamingDatasetReader.open forwards.

StreamingDatasetReader

StreamingDatasetReader(dataset: Any)

Thin facade over StreamingLeRobotDataset: iterate frames, or wrap in a DataLoader.

Build one with :meth:open (or the :func:stream_dataset alias). The underlying lerobot dataset is .dataset.

fps property

fps: Any

Frames per second the dataset was recorded at.

meta property

meta: Any

The dataset's LeRobotDatasetMetadata.

num_episodes property

num_episodes: Any

Total episodes in the dataset.

num_frames property

num_frames: Any

Total frames in the dataset.

dataloader

dataloader(batch_size: int = 64, num_workers: int = 0, **kw: Any) -> Any

Wrap the stream in a torch.utils.data.DataLoader.

shuffle is ignored: the stream shuffles internally (see :meth:open's Ordering note). With num_workers > 0 video decode parallelizes across the worker processes and must not ALSO run in the main process - lerobot documents a segfault when a second num_workers=0 loader touches the same video reader. lerobot's own make_dataset couples max_num_shards = num_workers; to match it, pass the same N to :meth:open's max_num_shards and here.

open classmethod

open(repo_id: str, *, root: str | None = None, episodes: list[int] | None = None, delta_timestamps: dict[str, list[float]] | None = None, image_transforms: Callable | None = None, tolerance_s: float = 0.0001, revision: str | None = None, streaming: bool = True, buffer_size: int = 1000, max_num_shards: int = 16, seed: int = 42, shuffle: bool = True, return_uint8: bool = True, validate_deltas: bool = True, drop_videos: bool = False, repo_type: str = 'dataset') -> StreamingDatasetReader

Open repo_id for streaming.

Parameters:

Name Type Description Default
repo_id str

Hub dataset id (org/name). A local directory a recorder registered under that id is used as root when root is not given.

required
root str | None

Local dataset directory; overrides the Hub.

None
episodes list[int] | None

Episode indices to stream (distinct, non-negative); all when None. lerobot's streaming dataset stores the list without reading it, so the subset is applied on iteration here; an index the dataset does not hold is refused rather than streamed as an empty read.

None
delta_timestamps dict[str, list[float]] | None

Per-key time offsets (seconds) to stack per frame. Checked against the dataset's fps grid when validate_deltas.

None
image_transforms Callable | None

Callable applied to every image tensor.

None
tolerance_s float

Half-width of the grid-match window (>= 0).

0.0001
revision str | None

Hub revision (branch, tag or commit).

None
streaming bool

False materializes the dataset instead.

True
buffer_size int

Reservoir the reader yields from, in frames (> 0); 1 is half of capture order - see Ordering.

1000
max_num_shards int

Parquet shards interleaved (> 0); 1 is the other half of capture order - see Ordering.

16
seed int

Shuffle seed (>= 0).

42
shuffle bool

Whether the reorder is reproducible ACROSS exhaustions - NOT whether it happens - see Ordering.

True
return_uint8 bool

Stream images as uint8 (a quarter of float32's bandwidth); policies normalize either.

True
validate_deltas bool

Refuse delta_timestamps off the fps grid.

True
drop_videos bool

Proprio-only streaming with no video decode (no torchcodec needed). Requires delta_timestamps naming at least one non-video key; camera keys in it are dropped.

False
repo_type str

"dataset" (the versioned Hub namespace) or "bucket" (Hub storage buckets).

'dataset'
Ordering

shuffle=False alone still reads shuffled frames. StreamingLeRobotDataset reorders at two levels either way - it samples a shard at random per frame, then yields from a reservoir of buffer_size - and shuffle selects only which generator drives that reorder (one reseeded from seed on every exhaustion, versus the dataset's own advancing one), i.e. reproducibility ACROSS epochs. Capture order therefore needs buffer_size=1 (a reservoir of one has nothing to reorder) together with max_num_shards=1 (a single shard has nothing to interleave). Nothing reports a shuffled read, so an eval or replay loop that asked for order with shuffle alone silently consumes frames out of order.

Raises:

Type Description
ValueError

A flag that is not a boolean, a numeric knob outside its domain, an episodes list the dataset cannot satisfy, or drop_videos=True without a usable delta_timestamps.

ImportError

lerobot's streaming dataset is not importable.

stream_dataset

stream_dataset(repo_id: str, **kwargs: Any) -> StreamingDatasetReader

Open a streaming reader for a LeRobotDataset without a simulator.

Module-level alias for :meth:StreamingDatasetReader.open; every keyword is forwarded unchanged.

has_streaming_dataset

has_streaming_dataset() -> bool

Return True if lerobot's StreamingLeRobotDataset is importable.

Transfer

strands_robots.dataset_transfer.sync_dataset_to_bucket

sync_dataset_to_bucket(root: str | Path, bucket: str, run_id: str | None = None, *, create: bool = True, private: bool = True, delete: bool = False) -> dict[str, Any]

Sync an on-disk LeRobotDataset into an HF Storage Bucket (Phase 1/2).

Lifecycle-independent: needs only a finalized dataset directory on disk (meta/ present) and the hf CLI - no live :class:~strands_robots.dataset_recorder.DatasetRecorder, no sim world. Covers syncing a dataset recorded earlier in the process, one recorded on hardware via lerobot-record, or a daily re-sync of a directory that grew. Both :meth:~strands_robots.dataset_recorder.DatasetRecorder.sync_to_bucket and the idle-path bucket sync in stop_recording delegate here so input validation and CLI orchestration exist exactly once.

Mutable, Xet-deduplicated dump target for COLLECTION - avoids git-LFS history bloat of push_to_hub during recording. Daily re-sync uploads only changed chunks (content-defined chunking). Requires the hf CLI with the buckets/sync subcommands (huggingface_hub>=1.5) and hf auth login.

bucket and run_id are validated against an allowlist before any subprocess or URI interpolation: bucket must be "name" or "org/name" and run_id a single path segment, both restricted to [A-Za-z0-9._-] (no path traversal or shell metacharacters). This path is agent-reachable via stop_recording(bucket=, run_id=). A rejected value returns {"status": "error", ...} without running hf.

create, private and delete select postures rather than scaling a quantity, so each is checked against :func:~strands_robots.utils.boolean_flag_error before the hf CLI is even located - the same domain the mesh provisioning entry points apply to their own capability flags. Read by truthiness they fail toward the permissive posture in both directions, because every non-empty string is truthy and every falsy non-boolean takes the other branch: delete="false" - the spelling an operator reaches for when opting out - appends --delete and mirror-deletes remote files absent locally, while private=0 drops --private and creates the bucket public.

The shard layout is already Xet/bucket-friendly at lerobot's defaults (100 MB data parquet / 200 MB video MP4 shards), and meta/ MUST ship or downstream loses normalization stats.

Parameters:

Name Type Description Default
root str | Path

Local dataset directory, str or Path (must contain meta/).

required
bucket str

Bucket target, "name" or "org/name".

required
run_id str | None

Subpath inside the bucket; defaults to the dataset directory name (Path(root).name).

None
create bool

Create the bucket first (pre-existing bucket is not an error). Must be a boolean.

True
private bool

Create the bucket as private (only used with create=True). Must be a boolean.

True
delete bool

Forward --delete to hf sync (mirror semantics - remove remote files absent locally). Must be a boolean.

False

Returns:

Type Description
dict[str, Any]

{"status": "success", "bucket_uri": ...} or

dict[str, Any]

{"status": "error", "message": ...}. Never raises on hf

dict[str, Any]

failure; errors are surfaced in the result dict. A flag outside its

dict[str, Any]

domain is reported the same way, without locating or running the CLI.

Episode judging

Judge-agent tools for labeling recorded LeRobotDataset episodes.

The four tools the episode-judge agent drives (see issue-level design in :mod:strands_robots.episode_labels, which owns the sidecar schema and the verdict-precedence contract):

  • :func:load_episode - episode/dataset metadata (frame count, features, whether a deterministic verdict and a judge label exist yet).
  • :func:sample_frames - evenly spaced frames from one episode: state vectors always, decoded camera images optionally (for a multimodal judge).
  • :func:read_predicate_verdict - the authoritative deterministic verdict.
  • :func:write_label - the judge's annotation. Structurally unable to overturn the deterministic verdict: it writes only the judge block, and a disagreeing success_opinion is recorded as disputes_verdict.

:func:create_judge_agent assembles them into a strands Agent whose system prompt carries the two-stage doctrine. It is model-provider agnostic: pass any strands model object (a Bedrock VLM, an OpenAI-compatible local endpoint, ...) or none for the strands default.

Every tool returns the structured {"status", "content"} envelope and never raises - a judge run over a hundred episodes must report the one episode it could not read, not die on it.

create_judge_agent

create_judge_agent(model: Any = None, system_prompt: str | None = None) -> Agent

Assemble the episode-judge agent from the four labeling tools.

Model-provider agnostic: model is any strands model object - a Bedrock-hosted VLM, an OpenAI-compatible local endpoint (vLLM/Ollama), or None for the strands default provider. No cloud dependency is required by the tools themselves; they only read the dataset and write the sidecar.

Parameters:

Name Type Description Default
model Any

Optional strands model object forwarded to Agent(model=...).

None
system_prompt str | None

Override for :data:JUDGE_SYSTEM_PROMPT. The default carries the two-stage doctrine (deterministic verdict is authoritative; the judge annotates).

None

Returns:

Type Description
Agent

A strands Agent wired with load_episode / sample_frames /

Agent

read_predicate_verdict / write_label.

load_episode

load_episode(root: str, episode: int) -> dict[str, Any]

Describe one recorded episode: length, features, and label state.

The judge's first call on an episode - it reports how many frames there are (bounds for sample_frames), which cameras and state columns the dataset carries, and whether a deterministic verdict / judge label already exists in the sidecar.

Parameters:

Name Type Description Default
root str

Dataset root directory (the directory containing meta/).

required
episode int

Episode index to describe, a non-negative whole number.

required

Returns:

Type Description
dict[str, Any]

Structured result whose JSON payload carries episode, length

dict[str, Any]

(frame count), total_episodes, fps, camera_keys,

dict[str, Any]

state_names, has_deterministic_verdict and has_judge_label.

sample_frames

sample_frames(root: str, episode: int, n_frames: int = 4, include_images: bool = False) -> dict[str, Any]

Sample evenly spaced frames from one episode for the judge to inspect.

Always returns the flattened observation.state vector and timestamp per sampled frame (pure pyarrow read, no lerobot import), plus a motion summary over the whole episode: rms_state_jerk (rms third difference of the state series) so a text-only judge can ground jerky_motion from state alone, and max_state_delta (largest per-step state delta) for spotting discontinuities and teleports. Smoothness lives in the jerk field - a maximum first difference is a peak-velocity statistic pinned by the gross traverse, so it cannot see a superimposed jitter. With include_images a multimodal judge additionally receives the decoded camera frames as image content blocks.

Parameters:

Name Type Description Default
root str

Dataset root directory (the directory containing meta/).

required
episode int

Episode index to sample, a non-negative whole number.

required
n_frames int

How many evenly spaced frames to sample, a positive count. Clamped to the episode length.

4
include_images bool

When True, decode the camera frames at the sampled positions into PNG image blocks (requires the lerobot extra; a dataset recorded without cameras reports an error). Must be a boolean - a posture flag is checked, never read by truthiness.

False

Returns:

Type Description
dict[str, Any]

Structured result whose JSON payload carries episode, length,

dict[str, Any]

samples (frame_index / timestamp / state per sample),

dict[str, Any]

max_state_delta and rms_state_jerk (state units per second

dict[str, Any]

cubed when timestamps are present, per step cubed otherwise; null

dict[str, Any]

when the episode is shorter than four frames). When images are

dict[str, Any]

requested, one image block follows per camera per sampled position -

dict[str, Any]

position-major, cameras in sorted key order within each position (the

dict[str, Any]

same order load_episode reports camera_keys) - the leading text

dict[str, Any]

block states the block count and that grouping, and every image block is

dict[str, Any]

immediately preceded by a text block naming its camera and its

dict[str, Any]

frame_index, so a judge reading a run of n_frames x n_cameras

dict[str, Any]

images can say which view a per-view observation belongs to and join

dict[str, Any]

that view back onto the state row for the same frame in samples. The label is adjacent to its

dict[str, Any]

image because that is what binds: the grouping sentence alone is a rule

dict[str, Any]

the judge must apply, and a rule stated at a distance from the images

dict[str, Any]

does not survive the flat run (measured on a three-camera recording with

dict[str, Any]

one view fully blocked, naming the blind camera scored 22/40 from the

dict[str, Any]

grouping sentence alone, 17/40 from the full block map spelled out in

dict[str, Any]

one text block, and 40/40 from these per-block labels, against 30/30 on

dict[str, Any]

the same frames asked one at a time). Every camera is deliberately

dict[str, Any]

included rather than one canonical view: the same world motion can be

dict[str, Any]

legible in one view and below a judge's threshold in another (measured

dict[str, Any]

on a real two-camera recording, where a 185 mm slide read as 84 px of

dict[str, Any]

travel in one view and 22 px in the other), so sampling a single

dict[str, Any]

camera would drop verdicts.

read_predicate_verdict

read_predicate_verdict(root: str, episode: int) -> dict[str, Any]

Read the authoritative deterministic predicate verdict for an episode.

The benchmark predicates scored the episode from simulator state at rollout time; that verdict is stage one of the two-stage labeling and the judge's annotation layers on top of it. If none is recorded the judge must stop: there is nothing to annotate.

Parameters:

Name Type Description Default
root str

Dataset root directory (the directory containing meta/).

required
episode int

Episode index, a non-negative whole number.

required

Returns:

Type Description
dict[str, Any]

Structured result whose JSON payload is the episode's

dict[str, Any]

deterministic block (success / failure and, when present,

dict[str, Any]

steps / cumulative_reward / seed).

write_label

write_label(root: str, episode: int, quality: str, failure_mode: str | None = None, note: str = '', success_opinion: bool | None = None, judge_model: str = '') -> dict[str, Any]

Write the judge's annotation for an episode into the label sidecar.

Annotation only, by construction: this writes the judge block and cannot reach the deterministic one, so no value passed here changes the benchmark verdict. A success_opinion that contradicts the verdict is recorded as disputes_verdict: true for human review - the verdict stands. An episode with no recorded deterministic verdict is refused.

Parameters:

Name Type Description Default
root str

Dataset root directory (the directory containing meta/).

required
episode int

Episode index to label, a non-negative whole number.

required
quality str

Quality grade, one of low / medium / high. Grades the execution visible in the recording (smoothness, directness, control), not the outcome - the deterministic verdict already carries success/failure, so a clean failure can be medium or high and a jerky or lucky success can be low.

required
failure_mode str | None

Optional tag from the fixed taxonomy (jerky_motion, near_miss, camera_occlusion, wrong_but_lucky, drift, collision, incomplete, other). Legal on a successful episode too - near_miss and wrong_but_lucky are exactly the annotations that make a success worth excluding from training data.

None
note str

Short free-text observation backing the grade and tag.

''
success_opinion bool | None

The judge's own success read, or omit to offer none. Disagreement with the deterministic verdict is recorded as a dispute annotation, never applied.

None
judge_model str

Identifier of the labeling model (or "human"), stored for provenance and calibration.

''

Returns:

Type Description
dict[str, Any]

Structured result whose JSON payload is the updated episode record

dict[str, Any]

(episode_index / deterministic / judge).

Edit page