Skip to content

Layered

Revision

Bases: BaseModel

A node in the tree of layers.

Each revision adds one parquet layer on top of its parent's; the record it resolves to is that layer over its ancestors'. The layer location derives from id via layer_dir, it is not stored.

Holds the connection it was created/loaded with (con), so every method below reuses it without needing it passed at every call. It is a private attribute, not a field: it never round-trips through (de)serialization (e.g. across a Prefect process boundary) - a revision that comes back without one lazily reattaches the process-level default on next access.

Notes

resolver property

resolver: Resolver

This record's Resolver, built once from its ancestry.

The ancestry is truncated at the deepest materialised ancestor, so a deep tree resolves from few entries. Cached rather than rebuilt per call: layers are write-once, so the only thing that can change the truncation point is materialise, which clears this itself.

Notes

record property

record: Record

This revision's resolved view, as a Record.

The framework-agnostic view: narwhals frames, with no sign of how many layers were folded to produce them. resolver remains the DuckDB-shaped view, which datarecord.tools still builds from.

Notes

create classmethod

create(
    con: DuckDBPyConnection | None = None,
    parent: UUID | None = None,
) -> Self

Insert a new revision, letting the DB assign the UUID.

Source code in src/datarecord/layered/revision.py
@classmethod
def create(
    cls, con: DuckDBPyConnection | None = None, parent: UUID | None = None
) -> Self:
    """Insert a new revision, letting the DB assign the UUID."""
    con = con or default_connection()
    id_, parent_ = _insert(con, parent)
    record = cls(id=id_, parent=parent_)
    record._con = con
    return record

get classmethod

get(
    revision_id: UUID, con: DuckDBPyConnection | None = None
) -> Self

Load a revision by id.

Source code in src/datarecord/layered/revision.py
@classmethod
def get(cls, revision_id: UUID, con: DuckDBPyConnection | None = None) -> Self:
    """Load a revision by id."""
    con = con or default_connection()
    id_, parent = _fetch(con, revision_id)
    record = cls(id=id_, parent=parent)
    record._con = con
    return record

child

child() -> Self

Branch a new revision off this one.

Any node may be a parent: a layer is write-once, so a base cannot shift under its descendants.

Notes
Source code in src/datarecord/layered/revision.py
def child(self) -> Self:
    """Branch a new revision off this one.

    Any node may be a parent: a layer is write-once, so a base cannot
    shift under its descendants.

    Notes
    -----
    - [a layer's data is write-once](https://energy-models.github.io/datarecord/design/layers/#a-layers-data-is-write-once)
    """
    return type(self).create(self.con, parent=self.id)

materialise

materialise() -> None

Write this node's caches - owner maps and resolved dims.

A policy rather than a lifecycle step, and purely additive: it changes no answer, only how many layers a descendant's read touches. Once these exist, a read stops here instead of walking further up.

Notes
Source code in src/datarecord/layered/revision.py
def materialise(self) -> None:
    """Write this node's caches - owner maps and resolved dims.

    A policy rather than a lifecycle step, and purely additive: it changes
    no answer, only how many layers a descendant's read touches. Once these
    exist, a read stops here instead of walking further up.

    Notes
    -----
    - [materialised node caches](https://energy-models.github.io/datarecord/design/layers/#materialised-node-caches)
    """
    resolve.materialise(self.id, self.resolver.sources, self.con)
    self._resolver = None

ancestry

ancestry() -> list[UUID]

Revision ids along the root->self path, root first.

Notes
Source code in src/datarecord/layered/revision.py
def ancestry(self) -> list[UUID]:
    """Revision ids along the root->self path, root first.

    Notes
    -----
    - [materialised node caches](https://energy-models.github.io/datarecord/design/layers/#materialised-node-caches)
    """
    return ancestry(self.con, self.id)

write_record

write_record(
    revision_id: UUID | None,
    source: LayerData | RecordLike,
    con: DuckDBPyConnection,
    *,
    uri: str | None = None,
) -> None

Write source as revision_id's layer, which must not exist yet.

An existing layer directory is an error rather than an overwrite or a merge, so a whole-record write can never half-replace what a record holds. Keys are looked up one at a time and each file written before the next is built, so a lazily-building source does one read per file rather than one per key up front.

Parameters:

Name Type Description Default
revision_id UUID | None

The record whose layer this is; layer_dir derives the path. None only together with uri, for a standalone record that belongs to no record.

required
uri str | None

Write here instead of at the revision's own layer - how a Directory commit target produces a record outside the layer tree.

None
source LayerData | RecordLike

The layer's contents: a LayerData - a StagedSource for a NewChild commit, a Resolver for a Directory one - or a framework's own RecordLike, wrapped in a thin adapter reading its Frames through the same enumerate-and-read pairs. Validated against its own schema before anything is written.

required
con DuckDBPyConnection

Connection to write through.

required

Raises:

Type Description
FileExistsError

If the layer directory already exists.

ValueError

If a long frame is missing a long-schema column, or the schema declares a key dim no frame carries - either would make the fold misresolve the layer.

Notes
Source code in src/datarecord/layered/write.py
def write_record(
    revision_id: UUID | None,
    source: LayerData | RecordLike,
    con: DuckDBPyConnection,
    *,
    uri: str | None = None,
) -> None:
    """Write `source` as `revision_id`'s layer, which must not exist yet.

    An existing layer directory is an error rather than an overwrite or a merge,
    so a whole-record write can never half-replace what a record holds. Keys are
    looked up one at a time and each file written before the next is built, so a
    lazily-building source does one read per file rather than one per key up
    front.

    Parameters
    ----------
    revision_id
        The record whose layer this is; `layer_dir` derives the path.
        `None` only together with `uri`, for a standalone record that belongs
        to no record.
    uri
        Write here instead of at the revision's own layer - how a `Directory`
        commit target produces a record outside the layer tree.
    source
        The layer's contents: a `LayerData` - a `StagedSource` for a `NewChild`
        commit, a `Resolver` for a `Directory` one - or a framework's own
        `RecordLike`, wrapped in a thin adapter reading its `Frames` through the
        same enumerate-and-read pairs. Validated against its own schema before
        anything is written.
    con
        Connection to write through.

    Raises
    ------
    FileExistsError
        If the layer directory already exists.
    ValueError
        If a long frame is missing a long-schema column, or the schema declares a key
        dim no frame carries - either would make the fold misresolve the layer.

    Notes
    -----
    - [the long schema](https://energy-models.github.io/datarecord/design/format/#the-long-schema)
    - [writing a whole record](https://energy-models.github.io/datarecord/design/writing/)
    - [committing](https://energy-models.github.io/datarecord/design/working-record/#committing)
    - [module layout](https://energy-models.github.io/datarecord/design/module-layout/)
    """
    if uri is None:
        if revision_id is None:
            msg = "write_record needs a revision_id or a uri"
            raise ValueError(msg)
        base = layer_dir(revision_id)
    else:
        base = uri if uri.endswith("/") else uri + "/"
    local = "://" not in base
    if local and Path(base).exists():
        msg = f"layer {base} already exists; write_record creates a new layer (https://energy-models.github.io/datarecord/design/writing/)"
        raise FileExistsError(msg)

    data = (
        source if isinstance(source, LayerData) else _RecordLikeAsLayerData(source, con)
    )
    schema = data.schema
    if uri is None:
        # One schema for the whole tree (https://energy-models.github.io/datarecord/design/schema/#one-schema-per-record). The first layer written
        # declares it; every later one is checked against it, so a layer
        # cannot quietly redefine what an attribute means.
        _reconcile_schema(schema, con)

    # Staged then renamed, so a frame that fails validation part-way through
    # leaves no layer rather than half of one (https://energy-models.github.io/datarecord/design/writing/). Validation happens as each
    # frame is built, since building it twice would defeat the laziness.
    staging = f"{base.rstrip('/')}.staging/" if local else base
    if local:
        Path(staging).mkdir(parents=True)
    try:
        # A layer holds only data: a layered record's one schema lives beside
        # `layers/`, not inside any of them (https://energy-models.github.io/datarecord/design/schema/#one-schema-per-record). A standalone directory *is*
        # one record, so there the schema belongs in the directory.
        if local and uri is not None:
            with open(staging + "manifest.json", "w") as fh:
                fh.write(schema.model_dump_json())
        kinds = [
            ("dims", data.axes(), data.axis, "dims"),
            ("entities", data.entity_types(), data.entity_type, "dims/entity_type"),
            ("groups", data.groups(), data.group, "groups"),
            ("attributes", data.attributes(), data.attribute, "inputs"),
        ]
        # `outputs/` only for a source carrying results, so a record with none
        # produces a layer without the directory rather than an empty one (https://energy-models.github.io/datarecord/design/writing/).
        output_names = data.attributes("outputs")
        if output_names:
            kinds.append(
                (
                    "outputs",
                    output_names,
                    lambda name: data.attribute(name, "outputs"),
                    "outputs",
                )
            )
        # Each type's names, to check record-wide uniqueness once every component
        # frame has been seen (https://energy-models.github.io/datarecord/design/format/#entity-is-unique-across-types).
        tagged: list[DuckDBPyRelation] = []
        for kind, keys, read, subdir in kinds:
            for key in keys:
                rel = read(
                    key
                )  # looked up exactly once (https://energy-models.github.io/datarecord/design/writing/)
                if rel is None:
                    continue
                _validate_frame(rel, kind, key, schema)
                if kind == "entities":
                    tagged.append(rel.project("entity", lit(key).alias("entity_type")))
                _write_frame(
                    rel,
                    f"{staging}{subdir}/{key}.parquet",
                    schema,
                    # A per-type member file is indexed by `entity` and holds
                    # one column per attribute; the type is the file it is in,
                    # and `dims/entity.parquet` is what carries it for every
                    # later reader. A column repeating it here would be a third
                    # copy that can disagree.
                    drop=("entity_type",) if kind == "entities" else (),
                )
        _require_unique(tagged, con)
    except BaseException:
        if local:
            shutil.rmtree(staging, ignore_errors=True)
        raise
    if local:
        Path(staging).rename(base.rstrip("/"))