Skip to content
Roboto
Esc
↑↓navigate↵open⌘Jpreview
On this page

roboto.experimental.topics.read_plan

Module Contents

PLAN_VERSION

roboto.experimental.topics.read_plan.PLAN_VERSION: int = 1#View Source

Contract version stamped on every plan.

ReadPlan validation refuses a plan whose version it does not recognize, so a consumer on an older contract fails at parse time instead of misreading a newer plan.

ReadPlan

class roboto.experimental.topics.read_plan.ReadPlan(/, **data)#View Source

Bases: pydantic.BaseModel

Resolves a read of one topic over a time window into the files to fetch and how to interpret them.

Parameters

data Any

Attributes

ReadPlan.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

ReadPlan.partitions

partitions tuple[ReadPlanPartition, ...] = () #

One entry per partition in the window, each its own fetch-and-interpret plan.

Ordered by when each file’s data begins. Partitions backed by the same file are contiguous in this tuple, ordered by where their data_range slice starts in that file, so they arrive in the order their data sits in it. This orders whole partitions, not rows.

ReadPlan.plan_version

plan_version int = 1 #

Contract version of this plan. Validation refuses a version this model does not recognize.

ReadPlan.projection

The output fields a consumer projects decoded rows to.

ReadPlan.schema_

schema_ ReadPlanSchemaRef | None = None #

The schema the plan reads under. Serializes as schema.

None when no partition of the topic lies in the window (within the request’s restrictions), and the plan then has no partitions. A plan that names a schema may also have no partitions.

ReadPlan.topic_id

topic_id str #

The topic this plan reads.

ReadPlan.window

window TimeWindow | None #

The time window the caller asked to read, which the plan was resolved over.

A consumer selects rows by each partition’s ReadPlanPartition.window alone, which can be narrower than this one.

None only when the request omitted a bound and either its restrictions select no timestamped data for the topic, or the bound it gave lies past the far end of that data (a start after the data ends, or an end before it begins); the plan then has no partitions.

ReadPlanFieldRef

class roboto.experimental.topics.read_plan.ReadPlanFieldRef(/, **data)#View Source

Bases: pydantic.BaseModel

A schema field named by both its id and its path components in the schema.

Parameters

data Any

Attributes

ReadPlanFieldRef.field_id

field_id str #

Identifier of the schema field.

ReadPlanFieldRef.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

ReadPlanFieldRef.path

The field’s path components within the schema, from the root to the field.

ReadPlanObjectRef

class roboto.experimental.topics.read_plan.ReadPlanObjectRef(/, **data)#View Source

Bases: pydantic.BaseModel

Points to the file backing a scan task. A consumer fetches the file’s bytes from it.

Parameters

data Any

Attributes

ReadPlanObjectRef.fs_node_id

fs_node_id str #

Identifier of the backing file. This id is stable, so a consumer can cache on it.

ReadPlanObjectRef.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

ReadPlanObjectRef.size_bytes

size_bytes int | None = None #

The source object’s size in bytes.

ReadPlanPartition

class roboto.experimental.topics.read_plan.ReadPlanPartition(/, **data)#View Source

Bases: pydantic.BaseModel

Everything needed to fetch and interpret one in-window partition’s bytes.

Parameters

data Any

Attributes

ReadPlanPartition.data_range

data_range roboto.domain.topics.record.DataRange | None = None #

The slice of the file this partition owns, as (start, end), or None for the whole file.

start is the first owned position and end is one past the last, like a Python slice. Positions are in the addressing native to the file the slice was declared against: stored-row positions counted from 0, or nanoseconds of media time for video. Any slice a consumer decodes is a range of stored rows.

A slice appears when one file packs several partitions’ data (a LeRobot v3 data file, for example). Each of those partitions carries its own time_offset_ns and may cover the same instants as the partitions beside it in the file, so time alone cannot separate its rows from theirs. Restrict each scan task’s decode to this slice before filtering by time, or the read returns the neighboring partitions’ rows too.

A scan task may name a transformed representation of the file the slice was declared against (see ReadPlanScanTask.transformations) rather than that file itself. The slice still applies unchanged: a partition that owns a slice is served only by representations holding the same content at the same positions.

ReadPlanPartition.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

ReadPlanPartition.scan_tasks

scan_tasks tuple[ReadPlanScanTask, ...] = () #

The files to read for this partition; empty when the partition has no readable data.

A partition may have several scan tasks, each covering a subtree of the schema or the whole schema; the read takes each field from the highest-precedence scan task whose subtree contains it.

ReadPlanPartition.time_offset_ns

time_offset_ns int #

Offset a consumer adds to each decoded row timestamp; the same for every row in the partition.

ReadPlanPartition.timestamp

Where this partition’s row timestamps come from.

ReadPlanPartition.topic_part_id

topic_part_id str #

Identifier of the partition.

ReadPlanPartition.window

window TimeWindow #

The window this partition’s rows are selected in, in absolute Unix-epoch nanoseconds.

A consumer keeps only the rows whose stored timestamp, once time_offset_ns is added, falls inside this window, inclusive on both ends. It is the only time filter a read applies.

It is the caller’s window intersected with the part of this partition’s backing file the read is scoped to, so it can be narrower than ReadPlan.window: a read scoped to a Session narrows it when the Session holds that file over only part of the file’s time span. The narrowing is stated per file, so every partition backed by one file carries the same window; ReadPlanPartition.data_range separates the partitions packed into one file.

ReadPlanProjection

class roboto.experimental.topics.read_plan.ReadPlanProjection(/, **data)#View Source

Bases: pydantic.BaseModel

The output fields the plan resolves rows to.

The projection takes exactly one of two forms: either every field in the schema (all is true, and the field list is left implicit so the plan need not enumerate a large schema) or an explicit fields list.

Parameters

data Any

Attributes

ReadPlanProjection.all

all bool = False #

True when the projection covers every field in the schema.

ReadPlanProjection.all_fields()

classmethod all_fields()#View Source

Return a projection covering every field in the schema.

Attributes

ReadPlanProjection.fields

fields tuple[ReadPlanFieldRef, ...] | None = None #

The resolved field set when the read is narrowed; None when all is true.

ReadPlanProjection.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

ReadPlanProjection.narrowed()

classmethod narrowed(fields)#View Source

Return a projection narrowed to an explicit field set.

Parameters

fields Iterable[ReadPlanFieldRef]

ReadPlanScanTask

class roboto.experimental.topics.read_plan.ReadPlanScanTask(/, **data)#View Source

Bases: pydantic.BaseModel

One file to open, with the topic to read from it and the format and transformations needed to interpret it.

Which of the representations satisfying the governing selector backs a scan task is service policy and may change between releases; only the selector’s hard-filter matching rule is contract.

Parameters

data Any

Attributes

ReadPlanScanTask.format

The format the bytes are stored in; selects the decoder a consumer applies.

ReadPlanScanTask.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

ReadPlanScanTask.object

The single file this scan task resolves to.

ReadPlanScanTask.precedence

precedence int #

Where two scan tasks’ subtrees contain the same field, the read takes it from the higher-precedence one.

ReadPlanScanTask.subtree

subtree ReadPlanFieldRef | None = None #

The field subtree this scan task covers; None covers the whole schema.

ReadPlanScanTask.topic_name

topic_name str | None = None #

Name of the topic whose records this scan task reads from its file.

A file can hold several topics (an MCAP file does when its channels carry different topic names), so a reader selects the topic’s records by this name. None when the plan does not name the topic; a reader then reads the file as holding a single topic.

ReadPlanScanTask.transformations

transformations tuple[str, ...] = () #

Transformations applied to produce this variant, in order; empty on the original.

ReadPlanSchemaRef

class roboto.experimental.topics.read_plan.ReadPlanSchemaRef(/, **data)#View Source

Bases: pydantic.BaseModel

Identifies the single topic schema the plan uses.

Parameters

data Any

Attributes

ReadPlanSchemaRef.checksum

checksum str #

Checksum of the schema’s content. A consumer can cache the schema by this value.

ReadPlanSchemaRef.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

ReadPlanSchemaRef.schema_id

schema_id str #

Identifier of the resolved topic schema.

ReadPlanTimestamp

class roboto.experimental.topics.read_plan.ReadPlanTimestamp(/, **data)#View Source

Bases: pydantic.BaseModel

Where a partition’s row timestamps come from.

Timestamps are either read out of a schema field (kind is "schema_field", and field names which one) or taken from the storage envelope (message log or publish time), in which case no schema field is involved and field is None.

Parameters

data Any

Attributes

ReadPlanTimestamp.field

field ReadPlanFieldRef | None = None #

The schema field timestamps are read from; set exactly when kind is "schema_field".

ReadPlanTimestamp.kind

How timestamps are sourced: from a schema field, or from the storage envelope.

ReadPlanTimestamp.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

ReadPlanTimestamp.unit

unit str | None = None #

Time unit of the designated field’s stored values (a TimeUnit value, e.g. "ms").

Only meaningful for a "schema_field" source, and only set when the schema declares the field’s unit. None when the schema does not record one; a consumer then treats non-self-describing values as nanoseconds, the unit the plan’s windows and offsets are recorded in. Envelope-derived timestamps (message log/publish time) are always assumed nanoseconds.

TimeWindow

class roboto.experimental.topics.read_plan.TimeWindow(/, **data)#View Source

Bases: pydantic.BaseModel

A closed time window in absolute nanoseconds since the Unix epoch; both bounds inclusive.

The same shape serves the window a plan resolves over and the window each partition’s rows are selected in.

Parameters

data Any

Attributes

TimeWindow.end

end int #

Inclusive upper bound, in nanoseconds.

TimeWindow.model_config

model_config #

Configuration for the model, should be a dictionary conforming to [ConfigDict][pydantic.config.ConfigDict].

TimeWindow.start

start int #

Inclusive lower bound, in nanoseconds.

Was this page helpful?