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 ¶
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.
dataloader ¶
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 ( |
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
|
delta_timestamps
|
dict[str, list[float]] | None
|
Per-key time offsets (seconds) to stack per frame.
Checked against the dataset's fps grid when |
None
|
image_transforms
|
Callable | None
|
Callable applied to every image tensor. |
None
|
tolerance_s
|
float
|
Half-width of the grid-match window ( |
0.0001
|
revision
|
str | None
|
Hub revision (branch, tag or commit). |
None
|
streaming
|
bool
|
|
True
|
buffer_size
|
int
|
Reservoir the reader yields from, in frames
( |
1000
|
max_num_shards
|
int
|
Parquet shards interleaved ( |
16
|
seed
|
int
|
Shuffle seed ( |
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 |
True
|
drop_videos
|
bool
|
Proprio-only streaming with no video decode (no
torchcodec needed). Requires |
False
|
repo_type
|
str
|
|
'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 |
ImportError
|
lerobot's streaming dataset is not importable. |
stream_dataset ¶
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 ¶
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, |
required |
bucket
|
str
|
Bucket target, |
required |
run_id
|
str | None
|
Subpath inside the bucket; defaults to the dataset directory
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 |
True
|
delete
|
bool
|
Forward |
False
|
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
|
dict[str, Any]
|
|
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 thejudgeblock, and a disagreeingsuccess_opinionis recorded asdisputes_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 ¶
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 |
None
|
system_prompt
|
str | None
|
Override for :data: |
None
|
Returns:
| Type | Description |
|---|---|
Agent
|
A strands |
Agent
|
|
load_episode ¶
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 |
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 |
dict[str, Any]
|
(frame count), |
dict[str, Any]
|
|
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 |
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 |
False
|
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
Structured result whose JSON payload carries |
dict[str, Any]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
cubed when timestamps are present, per step cubed otherwise; |
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 |
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]
|
|
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 |
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 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 |
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]
|
|
dict[str, Any]
|
|
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 |
required |
episode
|
int
|
Episode index to label, a non-negative whole number. |
required |
quality
|
str
|
Quality grade, one of |
required |
failure_mode
|
str | None
|
Optional tag from the fixed taxonomy ( |
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 |
''
|
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
Structured result whose JSON payload is the updated episode record |
dict[str, Any]
|
( |