Skip to content

graflo.architecture.schema.physical_keys

Translate document keys between logical and stored property names.

The ingestion pipeline produces documents keyed by logical property names. A database may store some of those properties under other names, recorded in DatabaseProfile.vertex_property_names and edge_specs[].property_names. :class:PhysicalKeys is the one place those maps are applied to documents: the writer translates every document and field list it hands a backend, and the schema-aware reads (Connection.graph_neighbors / traverse) translate their arguments in and their results back out.

With no stored names on the profile every method returns its input unchanged (the same object), so the translation costs nothing for the common case.

Attributes

Classes

PhysicalKeys

Logical <-> stored name translation for one resolved profile.

Source code in graflo/architecture/schema/physical_keys.py
class PhysicalKeys:
    """Logical <-> stored name translation for one resolved profile."""

    def __init__(self, profile: DatabaseProfile):
        self._vertex = {
            vertex: dict(names)
            for vertex, names in profile.vertex_property_names.items()
            if names
        }
        self._vertex_back = {
            vertex: _invert(names) for vertex, names in self._vertex.items()
        }
        self._edge: dict[EdgeId, dict[str, str]] = {
            spec.edge_id: dict(spec.property_names)
            for spec in profile.edge_specs
            if spec.purpose is None and spec.property_names
        }

    @property
    def active(self) -> bool:
        """Whether any property is stored under a different name."""
        return bool(self._vertex) or bool(self._edge)

    # -- vertices ---------------------------------------------------------

    def vertex_fields(self, vertex: str, fields: Iterable[str]) -> list[str]:
        """Stored names for logical *fields* of *vertex*."""
        names = self._vertex.get(vertex, {})
        return [names.get(field, field) for field in fields]

    def vertex_doc(self, vertex: str, doc: dict[str, Any]) -> dict[str, Any]:
        """*doc* keyed by stored names; *doc* itself when nothing is renamed."""
        return _rekey(doc, self._vertex.get(vertex, {}))

    def vertex_docs(
        self, vertex: str, docs: list[dict[str, Any]]
    ) -> list[dict[str, Any]]:
        if vertex not in self._vertex:
            return docs
        return [self.vertex_doc(vertex, doc) for doc in docs]

    def logical_vertex_doc(self, vertex: str, doc: dict[str, Any]) -> dict[str, Any]:
        """A document read back from the database, keyed by logical names."""
        return _rekey(doc, self._vertex_back.get(vertex, {}))

    # -- edges ------------------------------------------------------------

    def edge_fields(self, edge_id: EdgeId, fields: Iterable[str]) -> list[str]:
        """Stored names for logical property *fields* of the schema edge *edge_id*."""
        names = self._edge.get(edge_id, {})
        return [names.get(field, field) for field in fields]

    def edge_triples(
        self,
        edge_id: EdgeId,
        source: str,
        target: str,
        triples: list[Any],
    ) -> list[Any]:
        """``(source_doc, target_doc, weights)`` triples keyed by stored names.

        *edge_id* is the schema edge the triples belong to, which differs from
        the container key for an edge whose relation is extracted per document.
        """
        edge_names = self._edge.get(edge_id, {})
        if not edge_names and source not in self._vertex and target not in self._vertex:
            return triples
        out: list[Any] = []
        for triple in triples:
            source_doc, target_doc, *rest = triple
            weights = rest[0] if rest else {}
            out.append(
                (
                    self.vertex_doc(source, source_doc),
                    self.vertex_doc(target, target_doc),
                    _rekey(weights, edge_names),
                    *rest[1:],
                )
            )
        return out

    def logical_edge_triples(
        self,
        edge_id: EdgeId,
        source: str,
        target: str,
        triples: list[Any],
    ) -> list[Any]:
        """Triples read back from the database, keyed by logical names."""
        back = _invert(self._edge.get(edge_id, {}))
        if not back and source not in self._vertex and target not in self._vertex:
            return triples
        out: list[Any] = []
        for triple in triples:
            source_doc, target_doc, *rest = triple
            weights = rest[0] if rest else {}
            out.append(
                [
                    self.logical_vertex_doc(source, source_doc),
                    self.logical_vertex_doc(target, target_doc),
                    _rekey(weights, back) if isinstance(weights, dict) else weights,
                    *rest[1:],
                ]
            )
        return out

    # -- whole containers -------------------------------------------------

    def container(self, gc: GraphContainer, edge_config: EdgeConfig) -> GraphContainer:
        """*gc* with every vertex and edge document keyed by stored names.

        For backends that take a whole container (native bulk load). Edge keys
        whose schema edge cannot be found are passed through: the writer skips
        them the same way.
        """
        if not self.active:
            return gc
        vertices = {
            vertex: self.vertex_docs(vertex, docs)
            for vertex, docs in gc.vertices.items()
        }
        edges: dict[EdgeId, list[Any]] = {}
        for edge_id, triples in gc.edges.items():
            source, target, _relation = edge_id
            schema_id = schema_edge_id(edge_config, edge_id)
            edges[edge_id] = (
                self.edge_triples(schema_id, source, target, triples)
                if schema_id is not None
                else triples
            )
        return gc.model_copy(update={"vertices": vertices, "edges": edges})

    # -- schema-aware reads -----------------------------------------------

    def anchor_key(
        self, vertex: str, key: str | dict[str, Any]
    ) -> str | dict[str, Any]:
        """An anchor given as a field mapping, keyed by stored names; a raw id as-is."""
        return self.vertex_doc(vertex, key) if isinstance(key, dict) else key

    def edge_filter(self, filters: Any, edge_ids: Iterable[EdgeId]) -> Any:
        """An edge filter over the edges *edge_ids*, naming stored attributes.

        One filter is applied to every edge a walk follows, so a property it
        names must be stored under one name across them; edges sharing a
        relation always are.

        Raises:
            ValueError: if a property the filter names is stored under
                different names on different edges.
        """
        if filters is None or not self._edge:
            return filters
        names: dict[str, str] = {}
        conflicts: set[str] = set()
        for edge_id in edge_ids:
            for logical, stored in self._edge.get(edge_id, {}).items():
                if names.setdefault(logical, stored) != stored:
                    conflicts.add(logical)
        if not names:
            return filters
        expression = (
            filters
            if isinstance(filters, FilterExpression)
            else FilterExpression.from_dict(filters)
        )
        named = _filter_fields(expression)
        clash = sorted(conflicts & named)
        if clash:
            raise ValueError(
                f"Edge filter on {clash}: stored under different names on the "
                "edges this walk follows; restrict edge_types to one relation"
            )
        return expression.rename_fields(names)

    def logical_container(
        self, gc: GraphContainer, edge_config: EdgeConfig
    ) -> GraphContainer:
        """A container read from the database, keyed by logical names.

        Vertex documents are translated per type. Edge rows are dicts of edge
        attributes (plus the endpoint keys, which are left alone); a row read
        through a declared inverse name is translated with its stored edge's
        names.
        """
        if not self.active:
            return gc
        vertices = {
            vertex: (
                [self.logical_vertex_doc(vertex, doc) for doc in docs]
                if vertex in self._vertex_back
                else docs
            )
            for vertex, docs in gc.vertices.items()
        }
        edges: dict[EdgeId, list[Any]] = {}
        for edge_id, rows in gc.edges.items():
            stored_id = _stored_edge_id(edge_config, edge_id)
            back = _invert(self._edge.get(stored_id, {})) if stored_id else {}
            edges[edge_id] = (
                [_rekey(row, back) if isinstance(row, dict) else row for row in rows]
                if back
                else rows
            )
        return gc.model_copy(update={"vertices": vertices, "edges": edges})

Attributes

active property

Whether any property is stored under a different name.

Methods:

__init__(profile)
Source code in graflo/architecture/schema/physical_keys.py
def __init__(self, profile: DatabaseProfile):
    self._vertex = {
        vertex: dict(names)
        for vertex, names in profile.vertex_property_names.items()
        if names
    }
    self._vertex_back = {
        vertex: _invert(names) for vertex, names in self._vertex.items()
    }
    self._edge: dict[EdgeId, dict[str, str]] = {
        spec.edge_id: dict(spec.property_names)
        for spec in profile.edge_specs
        if spec.purpose is None and spec.property_names
    }
anchor_key(vertex, key)

An anchor given as a field mapping, keyed by stored names; a raw id as-is.

Source code in graflo/architecture/schema/physical_keys.py
def anchor_key(
    self, vertex: str, key: str | dict[str, Any]
) -> str | dict[str, Any]:
    """An anchor given as a field mapping, keyed by stored names; a raw id as-is."""
    return self.vertex_doc(vertex, key) if isinstance(key, dict) else key
container(gc, edge_config)

gc with every vertex and edge document keyed by stored names.

For backends that take a whole container (native bulk load). Edge keys whose schema edge cannot be found are passed through: the writer skips them the same way.

Source code in graflo/architecture/schema/physical_keys.py
def container(self, gc: GraphContainer, edge_config: EdgeConfig) -> GraphContainer:
    """*gc* with every vertex and edge document keyed by stored names.

    For backends that take a whole container (native bulk load). Edge keys
    whose schema edge cannot be found are passed through: the writer skips
    them the same way.
    """
    if not self.active:
        return gc
    vertices = {
        vertex: self.vertex_docs(vertex, docs)
        for vertex, docs in gc.vertices.items()
    }
    edges: dict[EdgeId, list[Any]] = {}
    for edge_id, triples in gc.edges.items():
        source, target, _relation = edge_id
        schema_id = schema_edge_id(edge_config, edge_id)
        edges[edge_id] = (
            self.edge_triples(schema_id, source, target, triples)
            if schema_id is not None
            else triples
        )
    return gc.model_copy(update={"vertices": vertices, "edges": edges})
edge_fields(edge_id, fields)

Stored names for logical property fields of the schema edge edge_id.

Source code in graflo/architecture/schema/physical_keys.py
def edge_fields(self, edge_id: EdgeId, fields: Iterable[str]) -> list[str]:
    """Stored names for logical property *fields* of the schema edge *edge_id*."""
    names = self._edge.get(edge_id, {})
    return [names.get(field, field) for field in fields]
edge_filter(filters, edge_ids)

An edge filter over the edges edge_ids, naming stored attributes.

One filter is applied to every edge a walk follows, so a property it names must be stored under one name across them; edges sharing a relation always are.

Raises:

Type Description
ValueError

if a property the filter names is stored under different names on different edges.

Source code in graflo/architecture/schema/physical_keys.py
def edge_filter(self, filters: Any, edge_ids: Iterable[EdgeId]) -> Any:
    """An edge filter over the edges *edge_ids*, naming stored attributes.

    One filter is applied to every edge a walk follows, so a property it
    names must be stored under one name across them; edges sharing a
    relation always are.

    Raises:
        ValueError: if a property the filter names is stored under
            different names on different edges.
    """
    if filters is None or not self._edge:
        return filters
    names: dict[str, str] = {}
    conflicts: set[str] = set()
    for edge_id in edge_ids:
        for logical, stored in self._edge.get(edge_id, {}).items():
            if names.setdefault(logical, stored) != stored:
                conflicts.add(logical)
    if not names:
        return filters
    expression = (
        filters
        if isinstance(filters, FilterExpression)
        else FilterExpression.from_dict(filters)
    )
    named = _filter_fields(expression)
    clash = sorted(conflicts & named)
    if clash:
        raise ValueError(
            f"Edge filter on {clash}: stored under different names on the "
            "edges this walk follows; restrict edge_types to one relation"
        )
    return expression.rename_fields(names)
edge_triples(edge_id, source, target, triples)

(source_doc, target_doc, weights) triples keyed by stored names.

edge_id is the schema edge the triples belong to, which differs from the container key for an edge whose relation is extracted per document.

Source code in graflo/architecture/schema/physical_keys.py
def edge_triples(
    self,
    edge_id: EdgeId,
    source: str,
    target: str,
    triples: list[Any],
) -> list[Any]:
    """``(source_doc, target_doc, weights)`` triples keyed by stored names.

    *edge_id* is the schema edge the triples belong to, which differs from
    the container key for an edge whose relation is extracted per document.
    """
    edge_names = self._edge.get(edge_id, {})
    if not edge_names and source not in self._vertex and target not in self._vertex:
        return triples
    out: list[Any] = []
    for triple in triples:
        source_doc, target_doc, *rest = triple
        weights = rest[0] if rest else {}
        out.append(
            (
                self.vertex_doc(source, source_doc),
                self.vertex_doc(target, target_doc),
                _rekey(weights, edge_names),
                *rest[1:],
            )
        )
    return out
logical_container(gc, edge_config)

A container read from the database, keyed by logical names.

Vertex documents are translated per type. Edge rows are dicts of edge attributes (plus the endpoint keys, which are left alone); a row read through a declared inverse name is translated with its stored edge's names.

Source code in graflo/architecture/schema/physical_keys.py
def logical_container(
    self, gc: GraphContainer, edge_config: EdgeConfig
) -> GraphContainer:
    """A container read from the database, keyed by logical names.

    Vertex documents are translated per type. Edge rows are dicts of edge
    attributes (plus the endpoint keys, which are left alone); a row read
    through a declared inverse name is translated with its stored edge's
    names.
    """
    if not self.active:
        return gc
    vertices = {
        vertex: (
            [self.logical_vertex_doc(vertex, doc) for doc in docs]
            if vertex in self._vertex_back
            else docs
        )
        for vertex, docs in gc.vertices.items()
    }
    edges: dict[EdgeId, list[Any]] = {}
    for edge_id, rows in gc.edges.items():
        stored_id = _stored_edge_id(edge_config, edge_id)
        back = _invert(self._edge.get(stored_id, {})) if stored_id else {}
        edges[edge_id] = (
            [_rekey(row, back) if isinstance(row, dict) else row for row in rows]
            if back
            else rows
        )
    return gc.model_copy(update={"vertices": vertices, "edges": edges})
logical_edge_triples(edge_id, source, target, triples)

Triples read back from the database, keyed by logical names.

Source code in graflo/architecture/schema/physical_keys.py
def logical_edge_triples(
    self,
    edge_id: EdgeId,
    source: str,
    target: str,
    triples: list[Any],
) -> list[Any]:
    """Triples read back from the database, keyed by logical names."""
    back = _invert(self._edge.get(edge_id, {}))
    if not back and source not in self._vertex and target not in self._vertex:
        return triples
    out: list[Any] = []
    for triple in triples:
        source_doc, target_doc, *rest = triple
        weights = rest[0] if rest else {}
        out.append(
            [
                self.logical_vertex_doc(source, source_doc),
                self.logical_vertex_doc(target, target_doc),
                _rekey(weights, back) if isinstance(weights, dict) else weights,
                *rest[1:],
            ]
        )
    return out
logical_vertex_doc(vertex, doc)

A document read back from the database, keyed by logical names.

Source code in graflo/architecture/schema/physical_keys.py
def logical_vertex_doc(self, vertex: str, doc: dict[str, Any]) -> dict[str, Any]:
    """A document read back from the database, keyed by logical names."""
    return _rekey(doc, self._vertex_back.get(vertex, {}))
vertex_doc(vertex, doc)

doc keyed by stored names; doc itself when nothing is renamed.

Source code in graflo/architecture/schema/physical_keys.py
def vertex_doc(self, vertex: str, doc: dict[str, Any]) -> dict[str, Any]:
    """*doc* keyed by stored names; *doc* itself when nothing is renamed."""
    return _rekey(doc, self._vertex.get(vertex, {}))
vertex_docs(vertex, docs)
Source code in graflo/architecture/schema/physical_keys.py
def vertex_docs(
    self, vertex: str, docs: list[dict[str, Any]]
) -> list[dict[str, Any]]:
    if vertex not in self._vertex:
        return docs
    return [self.vertex_doc(vertex, doc) for doc in docs]
vertex_fields(vertex, fields)

Stored names for logical fields of vertex.

Source code in graflo/architecture/schema/physical_keys.py
def vertex_fields(self, vertex: str, fields: Iterable[str]) -> list[str]:
    """Stored names for logical *fields* of *vertex*."""
    names = self._vertex.get(vertex, {})
    return [names.get(field, field) for field in fields]

Functions:

schema_edge_id(edge_config, edge_id)

The declared edge a container key belongs to.

An exact match first, then the relation=None edge that carries per-document relations.

Source code in graflo/architecture/schema/physical_keys.py
def schema_edge_id(edge_config: EdgeConfig, edge_id: EdgeId) -> EdgeId | None:
    """The declared edge a container key belongs to.

    An exact match first, then the ``relation=None`` edge that carries
    per-document relations.
    """
    if edge_id in edge_config:
        return edge_id
    null_id = (edge_id[0], edge_id[1], None)
    if null_id in edge_config:
        return null_id
    return None