roboto.experimental.topics.read_plan#

Module Contents#

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

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.

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

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)

model_config#

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

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

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

Ordered by where each file’s data begins. This orders whole partitions, not rows.

plan_version: int = 1#

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

projection: ReadPlanProjection#

The output fields a consumer projects decoded rows to.

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 session, when the request names one), and the plan then has no partitions. A plan that names a schema may also have no partitions.

topic_id: str#

The topic this plan reads.

window: TimeWindow#

The time window the plan resolves over.

class roboto.experimental.topics.read_plan.ReadPlanExtent(/, **data)#

Bases: pydantic.BaseModel

A partition’s time bounds, clipped to the plan window.

Parameters:

data (Any)

max: int#

Inclusive upper bound, in absolute Unix-epoch nanoseconds.

min: int#

Inclusive lower bound, in absolute Unix-epoch nanoseconds.

model_config#

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

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

Bases: pydantic.BaseModel

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

Parameters:

data (Any)

field_id: str#

Identifier of the schema field.

model_config#

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

path: roboto.domain.topics.record.FieldPath#

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

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

Bases: pydantic.BaseModel

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

Parameters:

data (Any)

fs_node_id: str#

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

model_config#

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

size_bytes: int | None = None#

The source object’s size in bytes.

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

Bases: pydantic.BaseModel

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

Parameters:

data (Any)

extent: ReadPlanExtent#

The partition’s time bounds, clipped to the window.

model_config#

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

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.

time_offset_ns: int#

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

timestamp: ReadPlanTimestamp#

Where this partition’s row timestamps come from.

topic_part_id: str#

Identifier of the partition.

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

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)

all: bool = False#

True when the projection covers every field in the schema.

classmethod all_fields()#

Return a projection covering every field in the schema.

Return type:

ReadPlanProjection

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

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

model_config#

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

classmethod narrowed(fields)#

Return a projection narrowed to an explicit field set.

Parameters:

fields (Iterable[ReadPlanFieldRef])

Return type:

ReadPlanProjection

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

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)

format: roboto.domain.topics.RepresentationStorageFormat#

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

model_config#

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

object: ReadPlanObjectRef#

The single file this scan task resolves to.

precedence: int#

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

subtree: ReadPlanFieldRef | None = None#

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

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.

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

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

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

Bases: pydantic.BaseModel

Identifies the single topic schema the plan uses.

Parameters:

data (Any)

checksum: str#

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

model_config#

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

schema_id: str#

Identifier of the resolved topic schema.

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

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)

field: ReadPlanFieldRef | None = None#

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

kind: roboto.domain.topics.TimelineSourceKind#

from a schema field, or from the storage envelope.

Type:

How timestamps are sourced

model_config#

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

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, matching how the plan’s extents are recorded. Envelope-derived timestamps (message log/publish time) are always assumed nanoseconds.

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

Bases: pydantic.BaseModel

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

Parameters:

data (Any)

end: int#

Inclusive upper bound, in nanoseconds.

model_config#

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

start: int#

Inclusive lower bound, in nanoseconds.