roboto.experimental.topics.read_plan#
Module Contents#
- roboto.experimental.topics.read_plan.PLAN_VERSION: int = 1#
Contract version stamped on every plan.
ReadPlanvalidation 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.BaseModelResolves 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.Nonewhen 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.BaseModelA 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.BaseModelA 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.BaseModelPoints 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.BaseModelEverything 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.BaseModelThe output fields the plan resolves rows to.
The projection takes exactly one of two forms: either every field in the schema (
allis true, and the field list is left implicit so the plan need not enumerate a large schema) or an explicitfieldslist.- 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:
- fields: tuple[ReadPlanFieldRef, ...] | None = None#
The resolved field set when the read is narrowed;
Nonewhenallis 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:
- class roboto.experimental.topics.read_plan.ReadPlanScanTask(/, **data)#
Bases:
pydantic.BaseModelOne 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;
Nonecovers 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.
Nonewhen 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.BaseModelIdentifies 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.BaseModelWhere a partition’s row timestamps come from.
Timestamps are either read out of a schema field (
kindis"schema_field", andfieldnames which one) or taken from the storage envelope (message log or publish time), in which case no schema field is involved andfieldisNone.- Parameters:
data (Any)
- field: ReadPlanFieldRef | None = None#
The schema field timestamps are read from; set exactly when
kindis"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
TimeUnitvalue, e.g."ms").Only meaningful for a
"schema_field"source, and only set when the schema declares the field’s unit.Nonewhen 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.BaseModelA 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.