Skip to content

graflo.architecture.pipeline

Pipeline runtime (execution). Declarations live in graflo.architecture.contract.

Modules:

Name Description
runtime

Pipeline runtime: actors, assembly, and executor.

Attributes

__all__ = ['Actor', 'ActorConstants', 'ActorExecutor', 'ActorInitContext', 'ActorWrapper', 'DescendActor', 'EdgeActor', 'TransformActor', 'VertexActor', 'VertexRouterActor'] module-attribute

Classes

Actor

Bases: ABC

Abstract base class for all actors in the system.

Source code in graflo/architecture/pipeline/runtime/actor/base.py
class Actor(ABC):
    """Abstract base class for all actors in the system."""

    @abstractmethod
    def __call__(
        self, ctx: ExtractionContext, lindex: LocationIndex, *nargs, **kwargs
    ) -> ExtractionContext:
        """Execute the actor's main processing logic."""

    def fetch_important_items(self) -> dict[str, object]:
        """Get a dictionary of important items for string representation."""
        return {}

    def finish_init(self, init_ctx: ActorInitContext) -> None:
        """Complete initialization of the actor."""

    def init_transforms(self, init_ctx: ActorInitContext) -> None:
        """Initialize transformations for the actor."""

    def count(self) -> int:
        """Get the count of items processed by this actor."""
        return 1

    def references_vertices(self) -> set[VertexName]:
        """Return vertex names this actor references."""
        return set()

    def _filter_items(self, items: dict[str, object]) -> dict[str, object]:
        """Filter out None and empty items."""
        return {k: v for k, v in items.items() if v is not None and v}

    def _stringify_items(self, items: dict[str, object]) -> dict[str, str]:
        """Convert items to string representation."""
        return {
            k: ", ".join(str(x) for x in v) if isinstance(v, (tuple, list)) else str(v)
            for k, v in items.items()
        }

    def _fetch_items_from_dict(self, keys: tuple[str, ...]) -> dict[str, object]:
        """Helper method to extract items from instance dict for string representation."""
        return {k: self.__dict__[k] for k in keys if k in self.__dict__}

    def __str__(self) -> str:
        """Get string representation of the actor."""
        d = self.fetch_important_items()
        d = self._filter_items(d)
        d = self._stringify_items(d)
        d_list = [[k, d[k]] for k in sorted(d)]
        d_list_b = [type(self).__name__] + [": ".join(x) for x in d_list]
        return "\n".join(d_list_b)

    __repr__ = __str__

    def fetch_actors(self, level: int, edges: list) -> tuple[int, type, str, list]:
        """Fetch actor information for tree representation."""
        return level, type(self), str(self), edges

Attributes

__repr__ = __str__ class-attribute instance-attribute

Methods:

__call__(ctx, lindex, *nargs, **kwargs) abstractmethod

Execute the actor's main processing logic.

Source code in graflo/architecture/pipeline/runtime/actor/base.py
@abstractmethod
def __call__(
    self, ctx: ExtractionContext, lindex: LocationIndex, *nargs, **kwargs
) -> ExtractionContext:
    """Execute the actor's main processing logic."""
__str__()

Get string representation of the actor.

Source code in graflo/architecture/pipeline/runtime/actor/base.py
def __str__(self) -> str:
    """Get string representation of the actor."""
    d = self.fetch_important_items()
    d = self._filter_items(d)
    d = self._stringify_items(d)
    d_list = [[k, d[k]] for k in sorted(d)]
    d_list_b = [type(self).__name__] + [": ".join(x) for x in d_list]
    return "\n".join(d_list_b)
count()

Get the count of items processed by this actor.

Source code in graflo/architecture/pipeline/runtime/actor/base.py
def count(self) -> int:
    """Get the count of items processed by this actor."""
    return 1
fetch_actors(level, edges)

Fetch actor information for tree representation.

Source code in graflo/architecture/pipeline/runtime/actor/base.py
def fetch_actors(self, level: int, edges: list) -> tuple[int, type, str, list]:
    """Fetch actor information for tree representation."""
    return level, type(self), str(self), edges
fetch_important_items()

Get a dictionary of important items for string representation.

Source code in graflo/architecture/pipeline/runtime/actor/base.py
def fetch_important_items(self) -> dict[str, object]:
    """Get a dictionary of important items for string representation."""
    return {}
finish_init(init_ctx)

Complete initialization of the actor.

Source code in graflo/architecture/pipeline/runtime/actor/base.py
def finish_init(self, init_ctx: ActorInitContext) -> None:
    """Complete initialization of the actor."""
init_transforms(init_ctx)

Initialize transformations for the actor.

Source code in graflo/architecture/pipeline/runtime/actor/base.py
def init_transforms(self, init_ctx: ActorInitContext) -> None:
    """Initialize transformations for the actor."""
references_vertices()

Return vertex names this actor references.

Source code in graflo/architecture/pipeline/runtime/actor/base.py
def references_vertices(self) -> set[VertexName]:
    """Return vertex names this actor references."""
    return set()

ActorConstants

Constants used throughout the actor system.

Source code in graflo/architecture/pipeline/runtime/actor/base.py
class ActorConstants:
    """Constants used throughout the actor system."""

    DESCEND_KEY: str = "key"
    DRESSING_TRANSFORMED_VALUE_KEY: str = "__value__"

Attributes

DESCEND_KEY = 'key' class-attribute instance-attribute
DRESSING_TRANSFORMED_VALUE_KEY = '__value__' class-attribute instance-attribute

ActorExecutor

Owns runtime extraction and assembly orchestration for an ActorWrapper.

Source code in graflo/architecture/pipeline/runtime/executor.py
class ActorExecutor:
    """Owns runtime extraction and assembly orchestration for an ActorWrapper."""

    def __init__(self, root: ActorWrapper):
        self.root = root

    def extract(self, doc: dict) -> ExtractionContext:
        extraction_ctx = ExtractionContext()
        return self.root(extraction_ctx, doc=doc)

    def assemble(
        self, extraction_ctx: ExtractionContext
    ) -> defaultdict[GraphEntity, list]:
        assembly_ctx = AssemblyContext.from_extraction(extraction_ctx)
        return self.root.assemble(assembly_ctx)

    def assemble_result(self, extraction_ctx: ExtractionContext) -> GraphAssemblyResult:
        return GraphAssemblyResult(entities=self.assemble(extraction_ctx))

Attributes

root = root instance-attribute

Methods:

__init__(root)
Source code in graflo/architecture/pipeline/runtime/executor.py
def __init__(self, root: ActorWrapper):
    self.root = root
assemble(extraction_ctx)
Source code in graflo/architecture/pipeline/runtime/executor.py
def assemble(
    self, extraction_ctx: ExtractionContext
) -> defaultdict[GraphEntity, list]:
    assembly_ctx = AssemblyContext.from_extraction(extraction_ctx)
    return self.root.assemble(assembly_ctx)
assemble_result(extraction_ctx)
Source code in graflo/architecture/pipeline/runtime/executor.py
def assemble_result(self, extraction_ctx: ExtractionContext) -> GraphAssemblyResult:
    return GraphAssemblyResult(entities=self.assemble(extraction_ctx))
extract(doc)
Source code in graflo/architecture/pipeline/runtime/executor.py
def extract(self, doc: dict) -> ExtractionContext:
    extraction_ctx = ExtractionContext()
    return self.root(extraction_ctx, doc=doc)

ActorInitContext dataclass

Typed initialization state shared across actor tree.

Source code in graflo/architecture/pipeline/runtime/actor/base.py
@dataclass(slots=True)
class ActorInitContext:
    """Typed initialization state shared across actor tree."""

    vertex_config: VertexConfig
    edge_config: EdgeConfig
    transforms: dict[str, ProtoTransform]
    edge_derivation: EdgeDerivationRegistry = field(
        default_factory=EdgeDerivationRegistry
    )
    allowed_vertex_names: set[VertexName] | None = None
    infer_edges: bool = True
    infer_edge_only: set[EdgeId] = field(default_factory=set)
    infer_edge_except: set[EdgeId] = field(default_factory=set)
    strict_references: bool = False
    fail_fast: bool = False
    tolerate_transform_errors: bool = True
    target_db_flavor: DBType | None = None
    #: What each accumulator role of the resource can hold; ``None`` for any class.
    role_reach: dict[str, frozenset[VertexName] | None] = field(default_factory=dict)

Attributes

allowed_vertex_names = None class-attribute instance-attribute
edge_config instance-attribute
edge_derivation = field(default_factory=EdgeDerivationRegistry) class-attribute instance-attribute
fail_fast = False class-attribute instance-attribute
infer_edge_except = field(default_factory=set) class-attribute instance-attribute
infer_edge_only = field(default_factory=set) class-attribute instance-attribute
infer_edges = True class-attribute instance-attribute
role_reach = field(default_factory=dict) class-attribute instance-attribute
strict_references = False class-attribute instance-attribute
target_db_flavor = None class-attribute instance-attribute
tolerate_transform_errors = True class-attribute instance-attribute
transforms instance-attribute
vertex_config instance-attribute

Methods:

__init__(vertex_config, edge_config, transforms, edge_derivation=EdgeDerivationRegistry(), allowed_vertex_names=None, infer_edges=True, infer_edge_only=set(), infer_edge_except=set(), strict_references=False, fail_fast=False, tolerate_transform_errors=True, target_db_flavor=None, role_reach=dict())

ActorWrapper

Wrapper class for managing actor instances.

Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
class ActorWrapper:
    """Wrapper class for managing actor instances."""

    def __init__(self, *args: Any, **kwargs: Any) -> None:
        config = parse_root_config(*args, **kwargs)
        w = ActorWrapper.from_config(config)
        self.actor = w.actor
        self.init_ctx = w.init_ctx

    @property
    def vertex_config(self) -> VertexConfig:
        return self.init_ctx.vertex_config

    @property
    def edge_config(self) -> EdgeConfig:
        return self.init_ctx.edge_config

    @property
    def infer_edges(self) -> bool:
        return self.init_ctx.infer_edges

    @property
    def infer_edge_only(self) -> set[EdgeId]:
        return self.init_ctx.infer_edge_only

    @property
    def infer_edge_except(self) -> set[EdgeId]:
        return self.init_ctx.infer_edge_except

    @property
    def target_db_flavor(self) -> DBType | None:
        return self.init_ctx.target_db_flavor

    def _inverse_pairs(self) -> dict[str, str]:
        """``{relation: declared inverse}`` of this resource's edge config, built once.

        Keyed by the config object so a wrapper re-initialized against another
        schema does not mirror through a stale table.
        """
        edge_config = self.init_ctx.edge_config
        cached = self.__dict__.get("_inverse_pairs_cache")
        if cached is None or cached[0] is not edge_config:
            cached = (edge_config, inverse_map(edge_config.inverses))
            self.__dict__["_inverse_pairs_cache"] = cached
        return cached[1]

    def init_transforms(self, init_ctx: ActorInitContext) -> None:
        self.init_ctx = init_ctx
        self.actor.init_transforms(init_ctx)

    def finish_init(self, init_ctx: ActorInitContext) -> None:
        self.init_ctx = init_ctx
        self.actor.init_transforms(init_ctx)
        self.actor.finish_init(init_ctx)
        # What inference may write: the edges declared, and those the steps
        # registered while initialising -- not ones a record registers mid-cast.
        self.__dict__["_inferable_edge_ids"] = frozenset(
            edge_id for edge_id, _ in init_ctx.edge_config.items()
        )

    def count(self) -> int:
        return self.actor.count()

    @classmethod
    def from_config(cls, config: ActorConfig) -> ActorWrapper:
        if isinstance(config, VertexActorConfig):
            actor = VertexActor.from_config(config)
        elif isinstance(config, TransformActorConfig):
            actor = TransformActor.from_config(config)
        elif isinstance(config, EdgeActorConfig):
            actor = EdgeActor.from_config(config)
        elif isinstance(config, DescendActorConfig):
            actor = DescendActor.from_config(config)
        elif isinstance(config, VertexRouterActorConfig):
            actor = VertexRouterActor.from_config(config)
        else:
            raise ValueError(
                f"Expected VertexActorConfig, TransformActorConfig, EdgeActorConfig, "
                f"DescendActorConfig, or VertexRouterActorConfig, got {type(config)}"
            )
        wrapper = cls.__new__(cls)
        wrapper.actor = actor
        wrapper.init_ctx = ActorInitContext(
            vertex_config=VertexConfig(vertices=[]),
            edge_config=EdgeConfig(),
            transforms={},
            allowed_vertex_names=None,
            infer_edges=True,
            infer_edge_only=set(),
            infer_edge_except=set(),
        )
        return wrapper

    @classmethod
    def _from_step(cls, step: dict[str, Any]) -> ActorWrapper:
        config = validate_actor_step(normalize_actor_step(step))
        return cls.from_config(config)

    def __call__(
        self,
        ctx: ExtractionContext,
        lindex: LocationIndex | None = None,
        *nargs: Any,
        **kwargs: Any,
    ) -> ExtractionContext:
        if lindex is None:
            lindex = LocationIndex()
        ctx = self.actor(ctx, lindex, *nargs, **kwargs)
        return ctx

    def assemble(
        self, ctx: ExtractionContext | AssemblyContext | ActionContext
    ) -> defaultdict[GraphEntity, list]:
        if isinstance(ctx, AssemblyContext):
            assembly_ctx = ctx
        else:
            assembly_ctx = AssemblyContext.from_extraction(ctx)
        # Synthetic identities must exist before edges are assembled and before
        # docs are deduplicated on their identity fields: a hash/funnel vertex
        # keys on ``id``, so an empty ``id`` gives edges no endpoint key and
        # fuse_doc_basis no basis (it would fold the batch into one doc).
        ensure_assigned_uuids_in_acc_vertex(assembly_ctx.acc_vertex, self.vertex_config)
        ensure_digest_identities_in_acc_vertex(
            assembly_ctx.acc_vertex, self.vertex_config
        )
        assemble_edges(
            ctx=assembly_ctx,
            vertex_config=self.vertex_config,
            edge_config=self.edge_config,
            infer_edges=self.infer_edges,
            infer_edge_only=self.infer_edge_only,
            infer_edge_except=self.infer_edge_except,
            target_db_flavor=self.target_db_flavor,
            edge_derivation=self.init_ctx.edge_derivation,
            inverse_pairs=self._inverse_pairs(),
            inferable=self.__dict__.get("_inferable_edge_ids"),
        )

        for vertex_name, dd in assembly_ctx.acc_vertex.items():
            for vertex_list in dd.values():
                # Lookup-only observations locate existing vertices for edge
                # endpoints; they must not become writes. They stay in
                # acc_vertex, which edge rendering reads, and are dropped here.
                writable = [x.vertex for x in vertex_list if not x.lookup_only]
                if not writable:
                    continue
                vertex_list_updated = fuse_doc_basis(
                    writable,
                    tuple(self.vertex_config.identity_fields(vertex_name)),
                )
                vertex_list_updated = pick_unique_dict(vertex_list_updated)
                assembly_ctx.acc_global[vertex_name] += vertex_list_updated

        assembly_ctx = add_blank_collections(assembly_ctx, self.vertex_config)

        if isinstance(ctx, ActionContext):
            ctx.acc_global = assembly_ctx.acc_global
            return ctx.acc_global
        return assembly_ctx.acc_global

    @classmethod
    def from_dict(cls, data: dict | list) -> ActorWrapper:
        if isinstance(data, list):
            return cls(*data)
        return cls(**data)

    def assemble_tree(
        self,
        fig_path: Path | str | None = None,
        output_format: str = "pdf",
        output_dpi: int | None = None,
    ):
        """Draw this pipeline's actor tree, or return it as a graph.

        Delegates to :func:`graflo.plot.plotter.assemble_tree`, which is the
        one implementation. ``graflo.plot`` sits above this layer, so the
        import is made here rather than at module scope; a missing plotting
        extra is reported rather than raised.

        Args:
            fig_path: Where to write the figure; ``None`` returns the graph.
            output_format: Figure format, when writing one.
            output_dpi: Raster resolution, for ``png``.

        Returns:
            ``networkx.MultiDiGraph | None``: the tree when *fig_path* is
            ``None``, otherwise ``None``.
        """
        import logging

        logger = logging.getLogger(__name__)
        try:
            from graflo.plot.plotter import assemble_tree
        except ImportError as exc:
            logger.error("not able to import the plotting stack: %s", exc)
            return None
        return assemble_tree(
            self,
            fig_path=fig_path,
            output_format=output_format,
            output_dpi=output_dpi,
        )

    def fetch_actors(self, level: int, edges: list) -> tuple[int, type, str, list]:
        return self.actor.fetch_actors(level, edges)

    def collect_actors(self) -> list[Actor]:
        actors = [self.actor]
        if isinstance(self.actor, DescendActor):
            for descendant in self.actor.descendants:
                actors.extend(descendant.collect_actors())
        return actors

    def find_descendants(
        self,
        predicate: Callable[[ActorWrapper], bool] | None = None,
        *,
        actor_type: type[Actor] | None = None,
        **attr_in: Any,
    ) -> list[ActorWrapper]:
        if predicate is None:

            def _predicate(w: ActorWrapper) -> bool:
                if actor_type is not None and not isinstance(w.actor, actor_type):
                    return False
                for attr, allowed in attr_in.items():
                    if allowed is None:
                        continue
                    val = getattr(w.actor, attr, None)
                    if val not in allowed:
                        return False
                return True

            predicate = _predicate

        result: list[ActorWrapper] = []
        if isinstance(self.actor, DescendActor):
            for d in self.actor.descendants:
                if predicate(d):
                    result.append(d)
                result.extend(d.find_descendants(predicate=predicate))
        return result

    def remove_descendants_if(self, predicate: Callable[[ActorWrapper], bool]) -> None:
        if isinstance(self.actor, DescendActor):
            for d in list(self.actor.descendants):
                d.remove_descendants_if(predicate=predicate)
            self.actor._descendants[:] = [
                d
                for d in self.actor.descendants
                if not predicate(d)
                and not (isinstance(d.actor, DescendActor) and d.count() == 0)
            ]

Attributes

actor = w.actor instance-attribute
edge_config property
infer_edge_except property
infer_edge_only property
infer_edges property
init_ctx = w.init_ctx instance-attribute
target_db_flavor property
vertex_config property

Methods:

__call__(ctx, lindex=None, *nargs, **kwargs)
Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
def __call__(
    self,
    ctx: ExtractionContext,
    lindex: LocationIndex | None = None,
    *nargs: Any,
    **kwargs: Any,
) -> ExtractionContext:
    if lindex is None:
        lindex = LocationIndex()
    ctx = self.actor(ctx, lindex, *nargs, **kwargs)
    return ctx
__init__(*args, **kwargs)
Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
def __init__(self, *args: Any, **kwargs: Any) -> None:
    config = parse_root_config(*args, **kwargs)
    w = ActorWrapper.from_config(config)
    self.actor = w.actor
    self.init_ctx = w.init_ctx
assemble(ctx)
Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
def assemble(
    self, ctx: ExtractionContext | AssemblyContext | ActionContext
) -> defaultdict[GraphEntity, list]:
    if isinstance(ctx, AssemblyContext):
        assembly_ctx = ctx
    else:
        assembly_ctx = AssemblyContext.from_extraction(ctx)
    # Synthetic identities must exist before edges are assembled and before
    # docs are deduplicated on their identity fields: a hash/funnel vertex
    # keys on ``id``, so an empty ``id`` gives edges no endpoint key and
    # fuse_doc_basis no basis (it would fold the batch into one doc).
    ensure_assigned_uuids_in_acc_vertex(assembly_ctx.acc_vertex, self.vertex_config)
    ensure_digest_identities_in_acc_vertex(
        assembly_ctx.acc_vertex, self.vertex_config
    )
    assemble_edges(
        ctx=assembly_ctx,
        vertex_config=self.vertex_config,
        edge_config=self.edge_config,
        infer_edges=self.infer_edges,
        infer_edge_only=self.infer_edge_only,
        infer_edge_except=self.infer_edge_except,
        target_db_flavor=self.target_db_flavor,
        edge_derivation=self.init_ctx.edge_derivation,
        inverse_pairs=self._inverse_pairs(),
        inferable=self.__dict__.get("_inferable_edge_ids"),
    )

    for vertex_name, dd in assembly_ctx.acc_vertex.items():
        for vertex_list in dd.values():
            # Lookup-only observations locate existing vertices for edge
            # endpoints; they must not become writes. They stay in
            # acc_vertex, which edge rendering reads, and are dropped here.
            writable = [x.vertex for x in vertex_list if not x.lookup_only]
            if not writable:
                continue
            vertex_list_updated = fuse_doc_basis(
                writable,
                tuple(self.vertex_config.identity_fields(vertex_name)),
            )
            vertex_list_updated = pick_unique_dict(vertex_list_updated)
            assembly_ctx.acc_global[vertex_name] += vertex_list_updated

    assembly_ctx = add_blank_collections(assembly_ctx, self.vertex_config)

    if isinstance(ctx, ActionContext):
        ctx.acc_global = assembly_ctx.acc_global
        return ctx.acc_global
    return assembly_ctx.acc_global
assemble_tree(fig_path=None, output_format='pdf', output_dpi=None)

Draw this pipeline's actor tree, or return it as a graph.

Delegates to :func:graflo.plot.plotter.assemble_tree, which is the one implementation. graflo.plot sits above this layer, so the import is made here rather than at module scope; a missing plotting extra is reported rather than raised.

Parameters:

Name Type Description Default
fig_path Path | str | None

Where to write the figure; None returns the graph.

None
output_format str

Figure format, when writing one.

'pdf'
output_dpi int | None

Raster resolution, for png.

None

Returns:

Type Description

networkx.MultiDiGraph | None: the tree when fig_path is

None, otherwise None.

Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
def assemble_tree(
    self,
    fig_path: Path | str | None = None,
    output_format: str = "pdf",
    output_dpi: int | None = None,
):
    """Draw this pipeline's actor tree, or return it as a graph.

    Delegates to :func:`graflo.plot.plotter.assemble_tree`, which is the
    one implementation. ``graflo.plot`` sits above this layer, so the
    import is made here rather than at module scope; a missing plotting
    extra is reported rather than raised.

    Args:
        fig_path: Where to write the figure; ``None`` returns the graph.
        output_format: Figure format, when writing one.
        output_dpi: Raster resolution, for ``png``.

    Returns:
        ``networkx.MultiDiGraph | None``: the tree when *fig_path* is
        ``None``, otherwise ``None``.
    """
    import logging

    logger = logging.getLogger(__name__)
    try:
        from graflo.plot.plotter import assemble_tree
    except ImportError as exc:
        logger.error("not able to import the plotting stack: %s", exc)
        return None
    return assemble_tree(
        self,
        fig_path=fig_path,
        output_format=output_format,
        output_dpi=output_dpi,
    )
collect_actors()
Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
def collect_actors(self) -> list[Actor]:
    actors = [self.actor]
    if isinstance(self.actor, DescendActor):
        for descendant in self.actor.descendants:
            actors.extend(descendant.collect_actors())
    return actors
count()
Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
def count(self) -> int:
    return self.actor.count()
fetch_actors(level, edges)
Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
def fetch_actors(self, level: int, edges: list) -> tuple[int, type, str, list]:
    return self.actor.fetch_actors(level, edges)
find_descendants(predicate=None, *, actor_type=None, **attr_in)
Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
def find_descendants(
    self,
    predicate: Callable[[ActorWrapper], bool] | None = None,
    *,
    actor_type: type[Actor] | None = None,
    **attr_in: Any,
) -> list[ActorWrapper]:
    if predicate is None:

        def _predicate(w: ActorWrapper) -> bool:
            if actor_type is not None and not isinstance(w.actor, actor_type):
                return False
            for attr, allowed in attr_in.items():
                if allowed is None:
                    continue
                val = getattr(w.actor, attr, None)
                if val not in allowed:
                    return False
            return True

        predicate = _predicate

    result: list[ActorWrapper] = []
    if isinstance(self.actor, DescendActor):
        for d in self.actor.descendants:
            if predicate(d):
                result.append(d)
            result.extend(d.find_descendants(predicate=predicate))
    return result
finish_init(init_ctx)
Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
def finish_init(self, init_ctx: ActorInitContext) -> None:
    self.init_ctx = init_ctx
    self.actor.init_transforms(init_ctx)
    self.actor.finish_init(init_ctx)
    # What inference may write: the edges declared, and those the steps
    # registered while initialising -- not ones a record registers mid-cast.
    self.__dict__["_inferable_edge_ids"] = frozenset(
        edge_id for edge_id, _ in init_ctx.edge_config.items()
    )
from_config(config) classmethod
Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
@classmethod
def from_config(cls, config: ActorConfig) -> ActorWrapper:
    if isinstance(config, VertexActorConfig):
        actor = VertexActor.from_config(config)
    elif isinstance(config, TransformActorConfig):
        actor = TransformActor.from_config(config)
    elif isinstance(config, EdgeActorConfig):
        actor = EdgeActor.from_config(config)
    elif isinstance(config, DescendActorConfig):
        actor = DescendActor.from_config(config)
    elif isinstance(config, VertexRouterActorConfig):
        actor = VertexRouterActor.from_config(config)
    else:
        raise ValueError(
            f"Expected VertexActorConfig, TransformActorConfig, EdgeActorConfig, "
            f"DescendActorConfig, or VertexRouterActorConfig, got {type(config)}"
        )
    wrapper = cls.__new__(cls)
    wrapper.actor = actor
    wrapper.init_ctx = ActorInitContext(
        vertex_config=VertexConfig(vertices=[]),
        edge_config=EdgeConfig(),
        transforms={},
        allowed_vertex_names=None,
        infer_edges=True,
        infer_edge_only=set(),
        infer_edge_except=set(),
    )
    return wrapper
from_dict(data) classmethod
Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
@classmethod
def from_dict(cls, data: dict | list) -> ActorWrapper:
    if isinstance(data, list):
        return cls(*data)
    return cls(**data)
init_transforms(init_ctx)
Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
def init_transforms(self, init_ctx: ActorInitContext) -> None:
    self.init_ctx = init_ctx
    self.actor.init_transforms(init_ctx)
remove_descendants_if(predicate)
Source code in graflo/architecture/pipeline/runtime/actor/wrapper.py
def remove_descendants_if(self, predicate: Callable[[ActorWrapper], bool]) -> None:
    if isinstance(self.actor, DescendActor):
        for d in list(self.actor.descendants):
            d.remove_descendants_if(predicate=predicate)
        self.actor._descendants[:] = [
            d
            for d in self.actor.descendants
            if not predicate(d)
            and not (isinstance(d.actor, DescendActor) and d.count() == 0)
        ]

DescendActor

Bases: Actor

Actor for processing hierarchical data structures.

Source code in graflo/architecture/pipeline/runtime/actor/descend.py
class DescendActor(Actor):
    """Actor for processing hierarchical data structures."""

    def __init__(
        self,
        key: str | None,
        any_key: bool = False,
        *,
        _descendants: list[ActorWrapper] | None = None,
    ):
        self.key = key
        self.any_key = any_key
        self._descendants: list[ActorWrapper] = (
            list(_descendants) if _descendants else []
        )
        self._descendants_sorted = True
        self._descendants.sort(key=lambda x: _NodeTypePriority[type(x.actor)])

    def fetch_important_items(self) -> dict[str, Any]:
        items = self._fetch_items_from_dict(("key",))
        if self.any_key:
            items["any_key"] = True
        return items

    def add_descendant(self, d: ActorWrapper) -> None:
        self._descendants.append(d)
        self._descendants_sorted = False

    def count(self) -> int:
        return sum(d.count() for d in self.descendants)

    @property
    def descendants(self) -> list[ActorWrapper]:
        if not self._descendants_sorted:
            self._descendants.sort(key=lambda x: _NodeTypePriority[type(x.actor)])
            self._descendants_sorted = True
        return self._descendants

    @classmethod
    def from_config(cls, config: DescendActorConfig) -> DescendActor:
        from .wrapper import ActorWrapper

        wrappers = [ActorWrapper.from_config(c) for c in config.pipeline]
        return cls(key=config.key, any_key=config.any_key, _descendants=wrappers)

    def _infer_vertex_descendants_from_transforms(
        self, init_ctx: ActorInitContext
    ) -> None:
        from .transform import TransformActor
        from .vertex import VertexActor

        if any(isinstance(an.actor, VertexActor) for an in self.descendants):
            return

        transform_output_fields: set[str] = set()
        for an in self.descendants:
            if (
                isinstance(an.actor, TransformActor)
                and an.actor._rename_map is not None
            ):
                transform_output_fields.update(
                    str(v) for v in an.actor._rename_map.values()
                )

        if not transform_output_fields:
            return

        inferred_vertices: list[str] = []
        for vertex_name in sorted(init_ctx.vertex_config.vertex_set):
            identity_fields = {
                f for f in init_ctx.vertex_config.identity_fields(vertex_name)
            }
            if identity_fields and identity_fields.issubset(transform_output_fields):
                inferred_vertices.append(vertex_name)

        if not inferred_vertices:
            return

        existing_targets: set[str] = set()
        for an in self.descendants:
            existing_targets.update(
                str(v) for v in an.actor.references_vertices() if v is not None
            )
        for vertex_name in inferred_vertices:
            if vertex_name in existing_targets:
                continue
            from .wrapper import ActorWrapper

            self.add_descendant(
                ActorWrapper.from_config(VertexActorConfig(vertex=vertex_name))
            )
            logger.debug(
                "DescendActor: inferred implicit VertexActor(%s) from untargeted transform fields %s",
                vertex_name,
                sorted(transform_output_fields),
            )

    def init_transforms(self, init_ctx: ActorInitContext) -> None:
        for an in self.descendants:
            an.init_transforms(init_ctx)

    def finish_init(self, init_ctx: ActorInitContext) -> None:
        self.vertex_config = init_ctx.vertex_config
        self._infer_vertex_descendants_from_transforms(init_ctx)
        for an in self.descendants:
            an.finish_init(init_ctx)

    def _expand_document(self, doc: dict | list) -> list[tuple[str | None, Any]]:
        if self.key is not None:
            if isinstance(doc, dict) and self.key in doc:
                items = doc[self.key]
                aux = items if isinstance(items, list) else [items]
                return [(self.key, item) for item in aux]
            return []
        elif self.any_key:
            if isinstance(doc, dict):
                result = []
                for key, items in doc.items():
                    aux = items if isinstance(items, list) else [items]
                    result.extend([(key, item) for item in aux])
                return result
            return []
        else:
            if isinstance(doc, list):
                return [(None, item) for item in doc]
            return [(None, doc)]

    def __call__(self, ctx: Any, lindex: Any, *nargs: Any, **kwargs: Any) -> Any:
        doc: Any = kwargs.pop("doc")
        if doc is None:
            raise ValueError(f"{type(self).__name__}: doc should be provided")
        if not doc:
            return ctx

        doc_expanded = self._expand_document(doc)
        if not doc_expanded:
            return ctx

        logger.debug("Expanding %s items", len(doc_expanded))

        for idoc, (key, sub_doc) in enumerate(doc_expanded):
            logger.debug("Processing item %s/%s", idoc + 1, len(doc_expanded))
            extra_step = (idoc,) if key is None else (key, idoc)
            child_lindex = lindex.extend(extra_step)
            if isinstance(sub_doc, dict):
                nargs_tuple: tuple[Any, ...] = ()
            else:
                nargs_tuple = (sub_doc,)

            for j, anw in enumerate(self.descendants):
                logger.debug(
                    "%s: %s/%s",
                    type(anw.actor).__name__,
                    j + 1,
                    len(self.descendants),
                )
                if isinstance(sub_doc, dict) and isinstance(anw.actor, TransformActor):
                    buf = list(ctx.transform_buffer.get(child_lindex, []))
                    feed_doc = merge_observation_with_transform_buffer(sub_doc, buf)
                    child_kwargs = {**kwargs, "doc": feed_doc}
                elif isinstance(sub_doc, dict):
                    child_kwargs = {**kwargs, "doc": sub_doc}
                else:
                    child_kwargs = kwargs
                ctx = anw(ctx, child_lindex, *nargs_tuple, **child_kwargs)
        return ctx

    def fetch_actors(self, level: int, edges: list) -> tuple[int, type, str, list]:
        label_current = str(self)
        cname_current = type(self)
        hash_current = hash((level, cname_current, label_current))
        logger.info("%s, %s", hash_current, (level, cname_current, label_current))
        props_current = {"label": label_current, "class": cname_current, "level": level}
        for d in self.descendants:
            level_a, cname, label_a, edges_a = d.fetch_actors(level + 1, edges)
            hash_a = hash((level_a, cname, label_a))
            props_a = {"label": label_a, "class": cname, "level": level_a}
            edges = [(hash_current, hash_a, props_current, props_a)] + edges_a
        return level, type(self), str(self), edges

Attributes

any_key = any_key instance-attribute
descendants property
key = key instance-attribute

Methods:

__call__(ctx, lindex, *nargs, **kwargs)
Source code in graflo/architecture/pipeline/runtime/actor/descend.py
def __call__(self, ctx: Any, lindex: Any, *nargs: Any, **kwargs: Any) -> Any:
    doc: Any = kwargs.pop("doc")
    if doc is None:
        raise ValueError(f"{type(self).__name__}: doc should be provided")
    if not doc:
        return ctx

    doc_expanded = self._expand_document(doc)
    if not doc_expanded:
        return ctx

    logger.debug("Expanding %s items", len(doc_expanded))

    for idoc, (key, sub_doc) in enumerate(doc_expanded):
        logger.debug("Processing item %s/%s", idoc + 1, len(doc_expanded))
        extra_step = (idoc,) if key is None else (key, idoc)
        child_lindex = lindex.extend(extra_step)
        if isinstance(sub_doc, dict):
            nargs_tuple: tuple[Any, ...] = ()
        else:
            nargs_tuple = (sub_doc,)

        for j, anw in enumerate(self.descendants):
            logger.debug(
                "%s: %s/%s",
                type(anw.actor).__name__,
                j + 1,
                len(self.descendants),
            )
            if isinstance(sub_doc, dict) and isinstance(anw.actor, TransformActor):
                buf = list(ctx.transform_buffer.get(child_lindex, []))
                feed_doc = merge_observation_with_transform_buffer(sub_doc, buf)
                child_kwargs = {**kwargs, "doc": feed_doc}
            elif isinstance(sub_doc, dict):
                child_kwargs = {**kwargs, "doc": sub_doc}
            else:
                child_kwargs = kwargs
            ctx = anw(ctx, child_lindex, *nargs_tuple, **child_kwargs)
    return ctx
__init__(key, any_key=False, *, _descendants=None)
Source code in graflo/architecture/pipeline/runtime/actor/descend.py
def __init__(
    self,
    key: str | None,
    any_key: bool = False,
    *,
    _descendants: list[ActorWrapper] | None = None,
):
    self.key = key
    self.any_key = any_key
    self._descendants: list[ActorWrapper] = (
        list(_descendants) if _descendants else []
    )
    self._descendants_sorted = True
    self._descendants.sort(key=lambda x: _NodeTypePriority[type(x.actor)])
add_descendant(d)
Source code in graflo/architecture/pipeline/runtime/actor/descend.py
def add_descendant(self, d: ActorWrapper) -> None:
    self._descendants.append(d)
    self._descendants_sorted = False
count()
Source code in graflo/architecture/pipeline/runtime/actor/descend.py
def count(self) -> int:
    return sum(d.count() for d in self.descendants)
fetch_actors(level, edges)
Source code in graflo/architecture/pipeline/runtime/actor/descend.py
def fetch_actors(self, level: int, edges: list) -> tuple[int, type, str, list]:
    label_current = str(self)
    cname_current = type(self)
    hash_current = hash((level, cname_current, label_current))
    logger.info("%s, %s", hash_current, (level, cname_current, label_current))
    props_current = {"label": label_current, "class": cname_current, "level": level}
    for d in self.descendants:
        level_a, cname, label_a, edges_a = d.fetch_actors(level + 1, edges)
        hash_a = hash((level_a, cname, label_a))
        props_a = {"label": label_a, "class": cname, "level": level_a}
        edges = [(hash_current, hash_a, props_current, props_a)] + edges_a
    return level, type(self), str(self), edges
fetch_important_items()
Source code in graflo/architecture/pipeline/runtime/actor/descend.py
def fetch_important_items(self) -> dict[str, Any]:
    items = self._fetch_items_from_dict(("key",))
    if self.any_key:
        items["any_key"] = True
    return items
finish_init(init_ctx)
Source code in graflo/architecture/pipeline/runtime/actor/descend.py
def finish_init(self, init_ctx: ActorInitContext) -> None:
    self.vertex_config = init_ctx.vertex_config
    self._infer_vertex_descendants_from_transforms(init_ctx)
    for an in self.descendants:
        an.finish_init(init_ctx)
from_config(config) classmethod
Source code in graflo/architecture/pipeline/runtime/actor/descend.py
@classmethod
def from_config(cls, config: DescendActorConfig) -> DescendActor:
    from .wrapper import ActorWrapper

    wrappers = [ActorWrapper.from_config(c) for c in config.pipeline]
    return cls(key=config.key, any_key=config.any_key, _descendants=wrappers)
init_transforms(init_ctx)
Source code in graflo/architecture/pipeline/runtime/actor/descend.py
def init_transforms(self, init_ctx: ActorInitContext) -> None:
    for an in self.descendants:
        an.init_transforms(init_ctx)

EdgeActor

Bases: Actor

Actor for processing edge data.

Operates in three modes determined by configuration:

Static mode (from/to set): both vertex types are declared at config time. The schema Edge is created during finish_init and the __call__ path is unchanged from the original implementation.

Dynamic mode (at least one of source_role/target_role set, with source_type_field/target_type_field accepted as legacy aliases): vertex types for the dynamic side(s) are resolved at extraction time by looking up accumulator slots populated by an upstream VertexRouterActor (slot segment = role or type_field) or a VertexActor with a matching role. The schema Edge is created—or retrieved from cache—per unique (source_type, target_type, relation) triple encountered.

Multi-link mode (links list set): each item in links becomes a dedicated sub-EdgeActor that runs in sequence per document, emitting one edge intent each. Use when one flat document encodes multiple distinct relationships.

Source code in graflo/architecture/pipeline/runtime/actor/edge.py
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
class EdgeActor(Actor):
    """Actor for processing edge data.

    Operates in three modes determined by configuration:

    **Static mode** (``from``/``to`` set): both vertex types are declared at config
    time.  The schema ``Edge`` is created during ``finish_init`` and the
    ``__call__`` path is unchanged from the original implementation.

    **Dynamic mode** (at least one of ``source_role``/``target_role`` set, with
    ``source_type_field``/``target_type_field`` accepted as legacy aliases):
    vertex types for the dynamic side(s) are
    resolved at extraction time by looking up accumulator slots populated by an
    upstream ``VertexRouterActor`` (slot segment = ``role`` or ``type_field``) or a
    ``VertexActor`` with a matching ``role``.
    The schema ``Edge`` is created—or retrieved from cache—per unique
    ``(source_type, target_type, relation)`` triple encountered.

    **Multi-link mode** (``links`` list set): each item in ``links`` becomes a
    dedicated sub-``EdgeActor`` that runs in sequence per document, emitting one edge
    intent each.  Use when one flat document encodes multiple distinct relationships.
    """

    def __init__(self, config: EdgeActorConfig):
        # Multi-link mode: delegate each link to its own EdgeActor.
        if config.links:
            self._link_actors: list[EdgeActor] = [
                EdgeActor(_link_to_edge_actor_config(lk)) for lk in config.links
            ]
            # Null-out all single-intent state so the dispatch is unambiguous.
            self._source_slot_key = None
            self._target_slot_key = None
            self._static_source = None
            self._static_target = None
            self._relation_map: dict[str, str] = {}
            self._relation_map_only = False
            self._strict_edge_types = False
            self._rejected_edges: set[tuple[str, str, str | None]] = set()
            self._edge_cache: dict[tuple[str, str, str | None], Edge] = {}
            self._init_ctx: ActorInitContext | None = None
            self.derivation: EdgeDerivation = EdgeDerivation()
            self._pending_vertex_weights: list[Weight] = []
            self._static_relation = None
            self.edge: Edge | None = None
            self.vertex_config: VertexConfig | None = None
            self.edge_config: EdgeConfig | None = None
            self.allowed_vertex_names: set[VertexName] | None = None
            return

        self._link_actors = []

        self._source_slot_key = config.source_role
        self._target_slot_key = config.target_role
        # Static fallback for whichever side is not dynamic.
        self._static_source = config.source
        self._static_target = config.target
        self._relation_map = config.relation_map or {}
        self._relation_map_only = config.relation_map_only
        self._strict_edge_types = config.strict_edge_types
        self._rejected_edges = set()
        self._edge_cache = {}
        self._init_ctx = None

        self.derivation = config.derivation
        self._pending_vertex_weights = []

        # In dynamic/mixed mode the static relation (if set) is used as a fallback
        # when relation_field yields nothing.
        self._static_relation = None

        # Dynamic mode: at least one side is resolved at extraction time.
        # Static mode: both sides are fixed at config time.
        is_dynamic = (
            self._source_slot_key is not None or self._target_slot_key is not None
        )
        if not is_dynamic:
            payload: dict[str, Any] = {
                "source": config.source,
                "target": config.target,
            }
            if config.relation is not None:
                payload["relation"] = config.relation
            if config.description is not None:
                payload["description"] = config.description
            if config.properties:
                payload["properties"] = config.properties
            for item in config.vertex_weights:
                self._pending_vertex_weights.append(Weight.model_validate(item))
            self.edge: Edge | None = Edge.from_dict(payload)
        else:
            self.edge = None
            self._static_relation = config.relation

        self.vertex_config: VertexConfig | None = None
        self.edge_config: EdgeConfig | None = None
        self.allowed_vertex_names: set[VertexName] | None = None

    @property
    def relation_field(self) -> str | None:
        """Alias for tooling (e.g. plot labels)."""
        return self.derivation.relation_field

    @property
    def is_dynamic(self) -> bool:
        """Whether this step names its edges per document, and so registers them mid-cast."""
        if self._link_actors:
            return any(link.is_dynamic for link in self._link_actors)
        return self._source_slot_key is not None or self._target_slot_key is not None

    @classmethod
    def from_config(cls, config: EdgeActorConfig) -> EdgeActor:
        return cls(config)

    def fetch_important_items(self) -> dict[str, Any]:
        if self._link_actors:
            return {"links": str(len(self._link_actors))}
        items: dict[str, Any] = {}
        if self.edge is not None:
            items["source"] = self.edge.source
            items["target"] = self.edge.target
        else:
            if self._source_slot_key is not None:
                items["source_role"] = self._source_slot_key
            elif self._static_source is not None:
                items["source"] = self._static_source
            if self._target_slot_key is not None:
                items["target_role"] = self._target_slot_key
            elif self._static_target is not None:
                items["target"] = self._static_target
        for k in ("match_source", "match_target"):
            v = getattr(self.derivation, k)
            if v is not None:
                items[k] = v
        return items

    def finish_init(self, init_ctx: ActorInitContext) -> None:
        self._init_ctx = init_ctx
        self.vertex_config = init_ctx.vertex_config
        self.edge_config = init_ctx.edge_config
        self.allowed_vertex_names = init_ctx.allowed_vertex_names

        if self._link_actors:
            # Multi-link mode: delegate finish_init to each sub-actor.
            for la in self._link_actors:
                la.finish_init(init_ctx)
            return

        if init_ctx.strict_references:
            # Strict references close the schema: a relation found in the data
            # must be declared too.
            self._strict_edge_types = True

        if self.edge is not None:
            # Static mode: register schema Edge now.
            self._adopt_declared_relation()
            edge_id = self.edge.edge_id
            if init_ctx.strict_references:
                self._refuse_undeclared(edge_id)
            init_ctx.edge_config.update_edges(
                self.edge, vertex_config=self.vertex_config
            )
            if self.derivation.relation_from_key:
                init_ctx.edge_derivation.mark_relation_from_key(edge_id)
            if self._pending_vertex_weights:
                init_ctx.edge_derivation.merge_vertex_weights(
                    edge_id, self._pending_vertex_weights
                )
            self._register_endpoint_match(init_ctx, edge_id)
            self.edge = init_ctx.edge_config.edge_for(edge_id)
            self._check_inverse_emission(init_ctx, edge_id)
        else:
            # Dynamic mode: cache will be populated per-document.
            self._edge_cache.clear()
            self._register_endpoint_rule(
                init_ctx, self._static_source, self._static_target
            )
            self._check_inverse_emission(init_ctx, None)

    def _refuse_undeclared(self, edge_id: EdgeId) -> None:
        """Refuse a static step whose edge the schema does not declare.

        Checked before the step registers its edge, so the resource's edge
        config still holds only what the schema declares.
        """
        if self.edge_config is None or self.edge_config.declared(edge_id) is not None:
            return
        source, target, relation = edge_id
        declared = sorted(
            str(r) for r in self._declared_relations(source, target) if r is not None
        )
        raise ValueError(
            f"edge step {source} -> {target} (relation {relation!r}) writes an edge "
            "the schema does not declare"
            + (f"; declared between them: {declared}" if declared else "")
            + ". Declare it in edge_config, or name a declared relation"
        )

    def _declared_relations(self, source: str, target: str) -> set[str | None]:
        """Relations the edge config declares from *source* to *target*."""
        if self.edge_config is None:
            return set()
        return {
            edge.relation
            for edge in self.edge_config.edges
            if edge.source == source and edge.target == target
        }

    def _adopt_declared_relation(self) -> None:
        """Give a static step that names no relation the one its endpoints declare.

        Otherwise the step registers a relation-less edge beside the declared
        one, and the writer drops every edge it renders. Several declared
        relations cannot be chosen between, so the step is refused.
        """
        if (
            self.edge is None
            or self.edge.relation is not None
            or self._relation_from_data()
        ):
            return
        declared = self._declared_relations(self.edge.source, self.edge.target)
        if not declared or None in declared:
            return
        if len(declared) > 1:
            raise ValueError(
                f"edge step {self.edge.source} -> {self.edge.target} names no "
                f"relation, and {sorted(r for r in declared if r)} are declared "
                "between them; name one with `relation`"
            )
        (relation,) = declared
        payload = self.edge.to_dict(skip_defaults=True)
        payload["relation"] = relation
        self.edge = Edge.from_dict(payload)

    def _check_inverse_emission(
        self, init_ctx: ActorInitContext, edge_id: EdgeId | None
    ) -> None:
        """Tie ``emit_inverse`` to the declared inverses, as far as the step is known.

        A step whose endpoints and relation are all fixed names exactly one edge,
        so everything can be checked now: the relation must have a declared pair
        (a symmetric relation has no inverse edge -- its edges are undirected),
        and the inverse edge must be declared, because only a *materialized*
        inverse is written. A step whose relation or endpoints come from the data
        is checked per document at assembly instead.

        Raises:
            ValueError: naming the step's edge and what to declare.
        """
        if not self.derivation.emit_inverse:
            return
        derivation = self.derivation
        data_driven = (
            edge_id is None
            or edge_id[2] is None
            or derivation.relation_field is not None
            or derivation.relation_from_key
        )
        if data_driven:
            if derivation.uses_secondary_identity():
                # The writer learns which identity to match an endpoint on per
                # edge id, at load time; an inverse known only per document has
                # no id to register the swapped selectors under.
                raise ValueError(
                    "emit_inverse cannot be combined with source_match / "
                    "target_match on an edge step whose relation or endpoints come "
                    "from the data; write the inverse with a step of its own"
                )
            return

        assert edge_id is not None
        edge_config = init_ctx.edge_config
        refusal = inverse_emission_refusal(edge_config, edge_id)
        if refusal is not None:
            raise ValueError(f"emit_inverse on edge step {edge_id}: {refusal}")
        source, target, relation = edge_id
        inverse_id = (target, source, edge_config.inverse_of(relation))
        if inverse_id in edge_config and (
            derivation.uses_secondary_identity() or derivation.on_ambiguous is not None
        ):
            # The mirror's endpoints are this step's endpoints, swapped.
            init_ctx.edge_derivation.set_endpoint_match(
                inverse_id,
                EndpointMatch(
                    source=selector_for(derivation.target_match, target),
                    target=selector_for(derivation.source_match, source),
                    on_ambiguous=derivation.on_ambiguous,
                ),
            )

    def _register_endpoint_match(
        self, init_ctx: ActorInitContext, edge_id: EdgeId
    ) -> None:
        """Validate endpoint identity selectors and record them for the writer.

        Resolving here fails fast at manifest load with the vertex name and the
        declared alternatives, rather than mid-ingest on the first batch. A
        relation read from the data leaves the edge id unknown until then, so
        the selectors are recorded as a rule over every relation.
        """
        source_type, target_type, _ = edge_id
        self._validate_selectors(init_ctx, source_type, target_type)
        if not self._selects_endpoints():
            return
        if self._relation_from_data():
            self._add_endpoint_rule(init_ctx, source_type, target_type)
            return
        derivation = self.derivation
        init_ctx.edge_derivation.set_endpoint_match(
            edge_id,
            EndpointMatch(
                source=selector_for(derivation.source_match, source_type),
                target=selector_for(derivation.target_match, target_type),
                on_ambiguous=derivation.on_ambiguous,
            ),
        )

    def _register_endpoint_rule(
        self,
        init_ctx: ActorInitContext,
        source: str | None,
        target: str | None,
    ) -> None:
        """Validate and record the selectors of a step whose edges are named per document.

        *source* / *target* are the endpoints fixed at config time, ``None``
        for one a router role fills.
        """
        self._validate_selectors(init_ctx, source, target)
        if self._selects_endpoints():
            self._add_endpoint_rule(init_ctx, source, target)

    def _selects_endpoints(self) -> bool:
        """Whether the writer needs to know how this step matches its endpoints."""
        derivation = self.derivation
        return (
            derivation.uses_secondary_identity() or derivation.on_ambiguous is not None
        )

    def _relation_from_data(self) -> bool:
        derivation = self.derivation
        return derivation.relation_field is not None or derivation.relation_from_key

    def _add_endpoint_rule(
        self, init_ctx: ActorInitContext, source: str | None, target: str | None
    ) -> None:
        """Record a rule over the edges this step can write.

        The relation is fixed only when the step names one and reads none from
        the data.
        """
        derivation = self.derivation
        if self._relation_from_data():
            relation = None
        elif self.edge is not None:
            relation = self.edge.relation
        else:
            relation = self._static_relation
        init_ctx.edge_derivation.add_endpoint_rule(
            EndpointRule(
                source=source,
                target=target,
                relation=relation,
                source_match=derivation.source_match,
                target_match=derivation.target_match,
                on_ambiguous=derivation.on_ambiguous,
            )
        )

    def _validate_selectors(
        self, init_ctx: ActorInitContext, source: str | None, target: str | None
    ) -> None:
        roles = init_ctx.role_reach
        self._validate_selector(
            "source_match",
            self.derivation.source_match,
            source,
            _held(roles, self._source_slot_key),
        )
        self._validate_selector(
            "target_match",
            self.derivation.target_match,
            target,
            _held(roles, self._target_slot_key),
        )

    def _validate_selector(
        self,
        name: str,
        selector: EndpointSelector | None,
        endpoint: str | None,
        held: frozenset[str] | None,
    ) -> None:
        """Refuse a selector no class this endpoint can be declares.

        *endpoint* is the class fixed at config time, ``None`` for one a router
        role fills; *held* is what that role can hold, ``None`` for any class.
        A plain selector on such an endpoint is resolved against each routed
        class when the edge is rendered; a per-class one is checked here, entry
        by entry: on a fixed endpoint it may name only that class, on a role
        only the classes the role can hold.

        Raises:
            ValueError: naming the selector, the class and what it declares.
        """
        vertex_config = self.vertex_config
        if vertex_config is None:
            return
        if isinstance(selector, dict):
            unknown = sorted(set(selector) - vertex_config.vertex_set)
            if unknown:
                raise ValueError(
                    f"{name} names {unknown}, which are not among the classes "
                    "this resource can produce"
                )
            if endpoint is not None:
                foreign = sorted(set(selector) - {endpoint})
                if foreign:
                    raise ValueError(
                        f"{name} names {foreign}, but the endpoint is always "
                        f"{endpoint!r}"
                    )
            elif held is not None:
                foreign = sorted(set(selector) - held)
                if foreign:
                    raise ValueError(
                        f"{name} names {foreign}, which the role can never hold; "
                        f"it holds {sorted(held)}"
                    )
            for vertex, plain in selector.items():
                # Raises with the declared alternatives when a selector is unknown.
                vertex_config.match_fields(vertex, plain)
        elif endpoint is not None and selector not in (None, PRIMARY_IDENTITY_SELECTOR):
            vertex_config.match_fields(endpoint, selector)

    # ------------------------------------------------------------------
    # Dynamic-mode helpers
    # ------------------------------------------------------------------

    def _get_or_create_edge(
        self, source: str, target: str, relation: str | None
    ) -> Edge | None:
        key = (source, target, relation)
        cached = self._edge_cache.get(key)
        if cached is not None:
            return cached
        if key in self._rejected_edges:
            return None
        with _EDGE_REGISTRATION_LOCK:
            return self._create_edge_locked(key)

    def _create_edge_locked(self, key: tuple[str, str, str | None]) -> Edge | None:
        """Create and register the edge for *key*; the registration lock is held."""
        source, target, relation = key
        cached = self._edge_cache.get(key)
        if cached is not None:
            return cached
        # Skip what the schema does not declare, by the writer's rule: the edge
        # itself, or the relation-less template between its endpoints.
        if (
            self._strict_edge_types
            and self.edge_config is not None
            and self.edge_config.declared(key) is None
        ):
            if key not in self._rejected_edges:
                self._rejected_edges.add(key)
                logger.warning(
                    "Edge (%s, %s, %s) is skipped: the schema does not declare it "
                    "(reported once per step).",
                    source,
                    target,
                    relation,
                )
            return None
        edge = Edge(
            source=source,
            target=target,
            relation=relation,
            # Edges of one relation agree on `directed`; a relation named only
            # at ingest time follows the ones already declared.
            directed=(
                self.edge_config.directed_for(relation)
                if self.edge_config is not None
                else True
            ),
        )
        if self.vertex_config is not None:
            edge.finish_init(vertex_config=self.vertex_config)
        if self.edge_config is not None and self.vertex_config is not None:
            self.edge_config.update_edges(edge, vertex_config=self.vertex_config)
        self._edge_cache[key] = edge
        logger.debug(
            "EdgeActor: registered dynamic edge (%s, %s, %s)", source, target, relation
        )
        return edge

    def _find_type_at_slot(
        self, ctx: ExtractionContext, slot_lindex: LocationIndex
    ) -> str | None:
        """Scan acc_vertex to find which vertex type has data at *slot_lindex*."""
        for vtype, by_loc in ctx.acc_vertex.items():
            if by_loc.get(slot_lindex):
                return vtype
        return None

    # ------------------------------------------------------------------
    # Main dispatch
    # ------------------------------------------------------------------

    def __call__(
        self, ctx: ExtractionContext, lindex: LocationIndex, *nargs: Any, **kwargs: Any
    ) -> ExtractionContext:
        if self._link_actors:
            # Multi-link mode: run each sub-actor in sequence.
            for la in self._link_actors:
                ctx = la(ctx, lindex, *nargs, **kwargs)
            return ctx
        if self._source_slot_key is not None or self._target_slot_key is not None:
            return self._call_dynamic(ctx, lindex, **kwargs)
        return self._call_static(ctx, lindex, **kwargs)

    def _call_static(
        self, ctx: ExtractionContext, lindex: LocationIndex, **kwargs: Any
    ) -> ExtractionContext:
        """Static mode: unchanged behavior from original EdgeActor."""
        assert self.edge is not None
        if self.allowed_vertex_names is not None and (
            self.edge.source not in self.allowed_vertex_names
            or self.edge.target not in self.allowed_vertex_names
        ):
            return ctx
        if (
            self.allowed_vertex_names is None
            and self.vertex_config is not None
            and (
                self.edge.source not in self.vertex_config.vertex_set
                or self.edge.target not in self.vertex_config.vertex_set
            )
        ):
            return ctx

        der = None if self.derivation.is_empty() else self.derivation
        ctx.record_edge_intent(
            edge=self.edge,
            location=lindex,
            derivation=der,
        )
        return ctx

    def _call_dynamic(
        self, ctx: ExtractionContext, lindex: LocationIndex, **kwargs: Any
    ) -> ExtractionContext:
        """Dynamic / mixed mode: resolve dynamic side(s) from VRA accumulator slots.

        Source or target (but not both) may be statically declared; that side's
        type is taken directly from config rather than looked up in the accumulator.
        """
        raw_observation = kwargs.get("doc", {})
        if not isinstance(raw_observation, dict):
            logger.debug(
                "EdgeActor: expected dict observation, got %s, skipping",
                type(raw_observation).__name__,
            )
            return ctx

        buffer_items: list[Any] = list(ctx.transform_buffer.get(lindex, []))
        doc = merge_observation_with_transform_buffer(raw_observation, buffer_items)
        ctx.obs_buffer[lindex] = dict(doc)

        # --- source type ---
        if self._source_slot_key is not None:
            source_slot_lindex = lindex.extend((self._source_slot_key, 0))
            source_type = self._find_type_at_slot(ctx, source_slot_lindex)
            if source_type is None:
                logger.debug(
                    "EdgeActor: no vertex data at source slot '%s', skipping",
                    self._source_slot_key,
                )
                return ctx
        else:
            # Mixed mode: static source
            assert self._static_source is not None
            source_type = self._static_source

        if (
            self.vertex_config is not None
            and source_type not in self.vertex_config.vertex_set
        ):
            logger.debug(
                "EdgeActor: source type '%s' not in vertex_set, skipping", source_type
            )
            return ctx

        # --- target type ---
        if self._target_slot_key is not None:
            target_slot_lindex = lindex.extend((self._target_slot_key, 0))
            target_type = self._find_type_at_slot(ctx, target_slot_lindex)
            if target_type is None:
                logger.debug(
                    "EdgeActor: no vertex data at target slot '%s', skipping",
                    self._target_slot_key,
                )
                return ctx
        else:
            # Mixed mode: static target
            assert self._static_target is not None
            target_type = self._static_target

        if (
            self.vertex_config is not None
            and target_type not in self.vertex_config.vertex_set
        ):
            logger.debug(
                "EdgeActor: target type '%s' not in vertex_set, skipping", target_type
            )
            return ctx

        # allowed_vertex_names early-exit
        if self.allowed_vertex_names is not None and (
            source_type not in self.allowed_vertex_names
            or target_type not in self.allowed_vertex_names
        ):
            return ctx

        # --- relation ---
        raw_relation: str | None
        if self.derivation.relation_field:
            raw_relation = doc.get(self.derivation.relation_field)
        else:
            raw_relation = None

        if raw_relation is not None:
            if self._relation_map_only and raw_relation not in self._relation_map:
                return ctx
            relation: str | None = self._relation_map.get(raw_relation, raw_relation)
        else:
            relation = self._static_relation
        if relation is None and not self._relation_from_data():
            # The per-document counterpart of `_adopt_declared_relation`; with
            # several declared relations the edge stays relation-less.
            declared = self._declared_relations(source_type, target_type)
            if len(declared) == 1:
                (relation,) = declared

        # Create / retrieve cached schema Edge.
        edge = self._get_or_create_edge(source_type, target_type, relation)
        if edge is None:
            return ctx

        # Build derivation: slot names for dynamic sides so render_edge can filter,
        # and the endpoint selectors, which render reduces to the concrete classes.
        derivation = EdgeDerivation(
            match_source=self._source_slot_key,
            match_target=self._target_slot_key,
            source_match=self.derivation.source_match,
            target_match=self.derivation.target_match,
            on_ambiguous=self.derivation.on_ambiguous,
            emit_inverse=self.derivation.emit_inverse,
        )
        ctx.record_edge_intent(edge=edge, location=lindex, derivation=derivation)
        return ctx

    def references_vertices(self) -> set[VertexName]:
        if self._link_actors:
            result: set[str] = set()
            for la in self._link_actors:
                result |= la.references_vertices()
            return result
        if self.edge is not None:
            return {self.edge.source, self.edge.target}
        static: set[str] = set()
        if self._static_source:
            static.add(self._static_source)
        if self._static_target:
            static.add(self._static_target)
        return (
            static
            | {s for s, _, _ in self._edge_cache}
            | {t for _, t, _ in self._edge_cache}
        )

Attributes

allowed_vertex_names = None instance-attribute
derivation = config.derivation instance-attribute
edge = None instance-attribute
edge_config = None instance-attribute
is_dynamic property

Whether this step names its edges per document, and so registers them mid-cast.

relation_field property

Alias for tooling (e.g. plot labels).

vertex_config = None instance-attribute

Methods:

__call__(ctx, lindex, *nargs, **kwargs)
Source code in graflo/architecture/pipeline/runtime/actor/edge.py
def __call__(
    self, ctx: ExtractionContext, lindex: LocationIndex, *nargs: Any, **kwargs: Any
) -> ExtractionContext:
    if self._link_actors:
        # Multi-link mode: run each sub-actor in sequence.
        for la in self._link_actors:
            ctx = la(ctx, lindex, *nargs, **kwargs)
        return ctx
    if self._source_slot_key is not None or self._target_slot_key is not None:
        return self._call_dynamic(ctx, lindex, **kwargs)
    return self._call_static(ctx, lindex, **kwargs)
__init__(config)
Source code in graflo/architecture/pipeline/runtime/actor/edge.py
def __init__(self, config: EdgeActorConfig):
    # Multi-link mode: delegate each link to its own EdgeActor.
    if config.links:
        self._link_actors: list[EdgeActor] = [
            EdgeActor(_link_to_edge_actor_config(lk)) for lk in config.links
        ]
        # Null-out all single-intent state so the dispatch is unambiguous.
        self._source_slot_key = None
        self._target_slot_key = None
        self._static_source = None
        self._static_target = None
        self._relation_map: dict[str, str] = {}
        self._relation_map_only = False
        self._strict_edge_types = False
        self._rejected_edges: set[tuple[str, str, str | None]] = set()
        self._edge_cache: dict[tuple[str, str, str | None], Edge] = {}
        self._init_ctx: ActorInitContext | None = None
        self.derivation: EdgeDerivation = EdgeDerivation()
        self._pending_vertex_weights: list[Weight] = []
        self._static_relation = None
        self.edge: Edge | None = None
        self.vertex_config: VertexConfig | None = None
        self.edge_config: EdgeConfig | None = None
        self.allowed_vertex_names: set[VertexName] | None = None
        return

    self._link_actors = []

    self._source_slot_key = config.source_role
    self._target_slot_key = config.target_role
    # Static fallback for whichever side is not dynamic.
    self._static_source = config.source
    self._static_target = config.target
    self._relation_map = config.relation_map or {}
    self._relation_map_only = config.relation_map_only
    self._strict_edge_types = config.strict_edge_types
    self._rejected_edges = set()
    self._edge_cache = {}
    self._init_ctx = None

    self.derivation = config.derivation
    self._pending_vertex_weights = []

    # In dynamic/mixed mode the static relation (if set) is used as a fallback
    # when relation_field yields nothing.
    self._static_relation = None

    # Dynamic mode: at least one side is resolved at extraction time.
    # Static mode: both sides are fixed at config time.
    is_dynamic = (
        self._source_slot_key is not None or self._target_slot_key is not None
    )
    if not is_dynamic:
        payload: dict[str, Any] = {
            "source": config.source,
            "target": config.target,
        }
        if config.relation is not None:
            payload["relation"] = config.relation
        if config.description is not None:
            payload["description"] = config.description
        if config.properties:
            payload["properties"] = config.properties
        for item in config.vertex_weights:
            self._pending_vertex_weights.append(Weight.model_validate(item))
        self.edge: Edge | None = Edge.from_dict(payload)
    else:
        self.edge = None
        self._static_relation = config.relation

    self.vertex_config: VertexConfig | None = None
    self.edge_config: EdgeConfig | None = None
    self.allowed_vertex_names: set[VertexName] | None = None
fetch_important_items()
Source code in graflo/architecture/pipeline/runtime/actor/edge.py
def fetch_important_items(self) -> dict[str, Any]:
    if self._link_actors:
        return {"links": str(len(self._link_actors))}
    items: dict[str, Any] = {}
    if self.edge is not None:
        items["source"] = self.edge.source
        items["target"] = self.edge.target
    else:
        if self._source_slot_key is not None:
            items["source_role"] = self._source_slot_key
        elif self._static_source is not None:
            items["source"] = self._static_source
        if self._target_slot_key is not None:
            items["target_role"] = self._target_slot_key
        elif self._static_target is not None:
            items["target"] = self._static_target
    for k in ("match_source", "match_target"):
        v = getattr(self.derivation, k)
        if v is not None:
            items[k] = v
    return items
finish_init(init_ctx)
Source code in graflo/architecture/pipeline/runtime/actor/edge.py
def finish_init(self, init_ctx: ActorInitContext) -> None:
    self._init_ctx = init_ctx
    self.vertex_config = init_ctx.vertex_config
    self.edge_config = init_ctx.edge_config
    self.allowed_vertex_names = init_ctx.allowed_vertex_names

    if self._link_actors:
        # Multi-link mode: delegate finish_init to each sub-actor.
        for la in self._link_actors:
            la.finish_init(init_ctx)
        return

    if init_ctx.strict_references:
        # Strict references close the schema: a relation found in the data
        # must be declared too.
        self._strict_edge_types = True

    if self.edge is not None:
        # Static mode: register schema Edge now.
        self._adopt_declared_relation()
        edge_id = self.edge.edge_id
        if init_ctx.strict_references:
            self._refuse_undeclared(edge_id)
        init_ctx.edge_config.update_edges(
            self.edge, vertex_config=self.vertex_config
        )
        if self.derivation.relation_from_key:
            init_ctx.edge_derivation.mark_relation_from_key(edge_id)
        if self._pending_vertex_weights:
            init_ctx.edge_derivation.merge_vertex_weights(
                edge_id, self._pending_vertex_weights
            )
        self._register_endpoint_match(init_ctx, edge_id)
        self.edge = init_ctx.edge_config.edge_for(edge_id)
        self._check_inverse_emission(init_ctx, edge_id)
    else:
        # Dynamic mode: cache will be populated per-document.
        self._edge_cache.clear()
        self._register_endpoint_rule(
            init_ctx, self._static_source, self._static_target
        )
        self._check_inverse_emission(init_ctx, None)
from_config(config) classmethod
Source code in graflo/architecture/pipeline/runtime/actor/edge.py
@classmethod
def from_config(cls, config: EdgeActorConfig) -> EdgeActor:
    return cls(config)
references_vertices()
Source code in graflo/architecture/pipeline/runtime/actor/edge.py
def references_vertices(self) -> set[VertexName]:
    if self._link_actors:
        result: set[str] = set()
        for la in self._link_actors:
            result |= la.references_vertices()
        return result
    if self.edge is not None:
        return {self.edge.source, self.edge.target}
    static: set[str] = set()
    if self._static_source:
        static.add(self._static_source)
    if self._static_target:
        static.add(self._static_target)
    return (
        static
        | {s for s, _, _ in self._edge_cache}
        | {t for _, t, _ in self._edge_cache}
    )

TransformActor

Bases: Actor

Actor for applying transformations to data.

Source code in graflo/architecture/pipeline/runtime/actor/transform.py
class TransformActor(Actor):
    """Actor for applying transformations to data."""

    def __init__(self, config: TransformActorConfig):
        self.transforms: dict[str, ProtoTransform] = {}
        self.call_use: str | None = None
        self._call_config = None
        self._fail_fast = False
        self._tolerate_transform_errors = True
        self._declared_input_keys: frozenset[str] = frozenset()
        self._rename_map: dict[str, str] | None = None
        self._guard: TransformGuardConfig | None = config.when

        if config.rename is not None:
            self.t = Transform(rename=config.rename)
            self._rename_map = dict(config.rename)
            return

        if config.call is None:
            raise ValueError(
                "TransformActorConfig requires call when rename is absent."
            )

        call = config.call
        self._call_config = call
        self.call_use = call.use
        inline_target = (
            call.target
            if call.target is not None
            else "values"
            if call.use is None
            else None
        )
        transform_kwargs: dict[str, Any] = {
            "name": call.use,
            "params": call.params,
            "module": call.module,
            "foo": call.foo,
            "input": tuple(call.input) if call.input else (),
            "output": tuple(call.output) if call.output else (),
            "input_groups": (
                tuple(tuple(group) for group in call.input_groups)
                if call.input_groups
                else ()
            ),
            "output_groups": (
                tuple(tuple(group) for group in call.output_groups)
                if call.output_groups
                else ()
            ),
            "dress": call.dress,
            "strategy": call.strategy or "single",
        }
        if inline_target is not None:
            transform_kwargs["target"] = inline_target
        if call.use is None and call.keys is not None:
            transform_kwargs["keys"] = KeySelectionConfig.model_validate(
                call.keys.model_dump()
            )
        # When call.use references ingestion_model.transforms, defer strict
        # transform validation until finish_init can hydrate module/foo.
        if call.use is not None and call.module is None and call.foo is None:
            self.t = Transform(name=call.use)
            return
        self.t = Transform(**transform_kwargs)

    def _refresh_missing_key_guard(self, init_ctx: ActorInitContext) -> None:
        self._fail_fast = init_ctx.fail_fast
        self._tolerate_transform_errors = init_ctx.tolerate_transform_errors
        if self.t.target == "keys" or self.t.strategy == "all":
            self._declared_input_keys = frozenset()
            return
        if self._rename_map is not None and not self._fail_fast:
            self._declared_input_keys = frozenset()
            return
        required: set[str] = set(self.t.input)
        for group in self.t.input_groups:
            required.update(group)
        self._declared_input_keys = frozenset(required)

    @staticmethod
    def _extract_observation(nargs: tuple[Any, ...], **kwargs: Any) -> Any:
        """Return the observation slice: dict observation or scalar/list positional value."""
        if kwargs:
            observation: Any | None = kwargs.get("doc")
        elif nargs:
            observation = nargs[0]
        else:
            raise ValueError(f"{type(TransformActor).__name__}: doc should be provided")
        if observation is None:
            raise ValueError(f"{type(TransformActor).__name__}: doc should be provided")
        return observation

    @staticmethod
    def _observation_keys(observation: Any) -> frozenset[str]:
        if isinstance(observation, dict):
            return frozenset(str(k) for k in observation)
        return frozenset()

    def _missing_declared_keys(self, observation: Any) -> frozenset[str]:
        if not self._declared_input_keys or not isinstance(observation, dict):
            return frozenset()
        return self._declared_input_keys - self._observation_keys(observation)

    def _rename_removed_keys(self, observation: Any) -> frozenset[str]:
        if self._rename_map is None or not isinstance(observation, dict):
            return frozenset()
        return frozenset(src for src in self._rename_map if src in observation)

    def fetch_important_items(self) -> dict[str, Any]:
        items = self._fetch_items_from_dict(("transform",))
        items.update({"t.input": self.t.input, "t.output": self.t.output})
        return items

    @classmethod
    def from_config(cls, config: TransformActorConfig) -> TransformActor:
        return cls(config)

    def init_transforms(self, init_ctx: ActorInitContext) -> None:
        self.transforms = init_ctx.transforms

    def _merge_call_with_proto(self, call: Any, pt: ProtoTransform) -> dict[str, Any]:
        next_params = call.params if call.params else pt.params
        next_dress = call.dress if call.dress is not None else pt.dress
        next_target = call.target if call.target is not None else pt.target

        if next_target == "keys":
            if call.input or call.output or call.input_groups or call.output_groups:
                raise ValueError(
                    "call.input, call.output, call.input_groups, and call.output_groups "
                    "cannot be used when the effective transform target is keys "
                    "(from call.target or the named ingestion_model.transforms entry)."
                )
            if call.dress is not None:
                raise ValueError("call.dress is not supported when target='keys'.")
            if call.strategy is not None and call.strategy != "single":
                raise ValueError(
                    "call.strategy is not allowed when target='keys'; "
                    "key mode uses implicit per-key execution."
                )
            next_input: tuple[str, ...] = ()
            next_output: tuple[str, ...] = ()
            next_input_groups: tuple[tuple[str, ...], ...] = ()
            next_output_groups: tuple[tuple[str, ...], ...] = ()
        else:
            next_input_groups = (
                tuple(tuple(group) for group in call.input_groups)
                if call.input_groups
                else pt.input_groups
            )
            next_output_groups = (
                tuple(tuple(group) for group in call.output_groups)
                if call.output_groups
                else pt.output_groups
            )
            if next_input_groups:
                next_input = ()
                # Explicit grouped override should not inherit potentially
                # conflicting proto output/output_groups for a different shape.
                if call.input_groups:
                    next_output_groups = (
                        tuple(tuple(group) for group in call.output_groups)
                        if call.output_groups
                        else ()
                    )
                    next_output = tuple(call.output) if call.output else ()
                elif next_dress is not None:
                    next_output = (next_dress.key, next_dress.value)
                else:
                    next_output = tuple(call.output) if call.output else pt.output
            else:
                next_input = tuple(call.input) if call.input else pt.input
                if next_dress is not None:
                    next_output = (next_dress.key, next_dress.value)
                else:
                    next_output = tuple(call.output) if call.output else pt.output

        transform_kwargs: dict[str, Any] = {
            "dress": next_dress,
            "name": call.use,
            "module": pt.module,
            "foo": pt.foo,
            "params": next_params,
            "input": next_input,
            "output": next_output,
            "input_groups": next_input_groups,
            "output_groups": next_output_groups,
            "strategy": call.strategy or "single",
            "target": next_target,
        }
        if call.keys is not None:
            transform_kwargs["keys"] = KeySelectionConfig.model_validate(
                call.keys.model_dump()
            )
        else:
            transform_kwargs["keys"] = pt.keys.model_copy(deep=True)
        return transform_kwargs

    def finish_init(self, init_ctx: ActorInitContext) -> None:
        self.transforms = init_ctx.transforms
        if self.call_use is None or self.t._foo is not None:
            self._refresh_missing_key_guard(init_ctx)
            return
        if self._call_config is None:
            self._refresh_missing_key_guard(init_ctx)
            return
        pt = self.transforms.get(self.call_use, None)
        if pt is None:
            if init_ctx.strict_references:
                raise ValueError(
                    f"Transform '{self.call_use}' referenced by transform.call.use "
                    "was not found in ingestion_model.transforms."
                )
            self._refresh_missing_key_guard(init_ctx)
            return
        call = self._call_config
        transform_kwargs = self._merge_call_with_proto(call, pt)
        self.t = Transform(**transform_kwargs)
        self._refresh_missing_key_guard(init_ctx)

    def _format_transform_result(self, result: Any) -> TransformPayload:
        return TransformPayload.from_result(result)

    def _transform_label(self) -> str:
        if self.call_use:
            return self.call_use
        if self.t.foo and self.t.module:
            return f"{self.t.module}.{self.t.foo}"
        if self.t.foo:
            return self.t.foo
        if self.t.name:
            return self.t.name
        return type(self.t).__name__

    def _format_traceback(self, exc: BaseException) -> str:
        return "".join(traceback.format_exception(type(exc), exc, exc.__traceback__))

    def __call__(
        self, ctx: ExtractionContext, lindex: LocationIndex, *nargs: Any, **kwargs: Any
    ) -> ExtractionContext:
        logger.debug("transforms : %s %s", id(self.transforms), len(self.transforms))
        observation = self._extract_observation(nargs, **kwargs)
        # A failed guard is a skip, not an error: the step writes nothing, which
        # is the property a guarded derivation behind a router relies on.
        if self._guard is not None and not self._guard.passes(observation):
            return ctx
        missing = self._missing_declared_keys(observation)
        if missing:
            if self._fail_fast:
                raise TransformException(
                    f"Missing required input keys: {sorted(missing)}"
                )
            return ctx
        try:
            transform_result = self.t(observation)
        except Exception as exc:
            if not self._tolerate_transform_errors:
                raise
            nulled_fields = self.t.planned_output_field_names(
                observation if isinstance(observation, dict) else None
            )
            if nulled_fields:
                payload = TransformPayload(named={k: None for k in nulled_fields})
                ctx.transform_buffer[lindex].append(payload)
                ctx.record_transform_observation(location=lindex, payload=payload)
            ctx.record_transform_failure(
                location=lindex,
                transform_label=self._transform_label(),
                exc=exc,
                traceback_text=self._format_traceback(exc),
                nulled_fields=nulled_fields,
            )
            return ctx
        if self._rename_map is not None:
            base = TransformPayload.from_result(transform_result)
            _update_doc = TransformPayload(
                named=base.named,
                positional=base.positional,
                removed_keys=self._rename_removed_keys(observation),
            )
        else:
            _update_doc = self._format_transform_result(transform_result)
        ctx.transform_buffer[lindex].append(_update_doc)
        ctx.record_transform_observation(location=lindex, payload=_update_doc)
        return ctx

    def references_vertices(self) -> set[str]:
        return set()

Attributes

call_use = call.use instance-attribute
t = Transform(**transform_kwargs) instance-attribute
transforms = {} instance-attribute

Methods:

__call__(ctx, lindex, *nargs, **kwargs)
Source code in graflo/architecture/pipeline/runtime/actor/transform.py
def __call__(
    self, ctx: ExtractionContext, lindex: LocationIndex, *nargs: Any, **kwargs: Any
) -> ExtractionContext:
    logger.debug("transforms : %s %s", id(self.transforms), len(self.transforms))
    observation = self._extract_observation(nargs, **kwargs)
    # A failed guard is a skip, not an error: the step writes nothing, which
    # is the property a guarded derivation behind a router relies on.
    if self._guard is not None and not self._guard.passes(observation):
        return ctx
    missing = self._missing_declared_keys(observation)
    if missing:
        if self._fail_fast:
            raise TransformException(
                f"Missing required input keys: {sorted(missing)}"
            )
        return ctx
    try:
        transform_result = self.t(observation)
    except Exception as exc:
        if not self._tolerate_transform_errors:
            raise
        nulled_fields = self.t.planned_output_field_names(
            observation if isinstance(observation, dict) else None
        )
        if nulled_fields:
            payload = TransformPayload(named={k: None for k in nulled_fields})
            ctx.transform_buffer[lindex].append(payload)
            ctx.record_transform_observation(location=lindex, payload=payload)
        ctx.record_transform_failure(
            location=lindex,
            transform_label=self._transform_label(),
            exc=exc,
            traceback_text=self._format_traceback(exc),
            nulled_fields=nulled_fields,
        )
        return ctx
    if self._rename_map is not None:
        base = TransformPayload.from_result(transform_result)
        _update_doc = TransformPayload(
            named=base.named,
            positional=base.positional,
            removed_keys=self._rename_removed_keys(observation),
        )
    else:
        _update_doc = self._format_transform_result(transform_result)
    ctx.transform_buffer[lindex].append(_update_doc)
    ctx.record_transform_observation(location=lindex, payload=_update_doc)
    return ctx
__init__(config)
Source code in graflo/architecture/pipeline/runtime/actor/transform.py
def __init__(self, config: TransformActorConfig):
    self.transforms: dict[str, ProtoTransform] = {}
    self.call_use: str | None = None
    self._call_config = None
    self._fail_fast = False
    self._tolerate_transform_errors = True
    self._declared_input_keys: frozenset[str] = frozenset()
    self._rename_map: dict[str, str] | None = None
    self._guard: TransformGuardConfig | None = config.when

    if config.rename is not None:
        self.t = Transform(rename=config.rename)
        self._rename_map = dict(config.rename)
        return

    if config.call is None:
        raise ValueError(
            "TransformActorConfig requires call when rename is absent."
        )

    call = config.call
    self._call_config = call
    self.call_use = call.use
    inline_target = (
        call.target
        if call.target is not None
        else "values"
        if call.use is None
        else None
    )
    transform_kwargs: dict[str, Any] = {
        "name": call.use,
        "params": call.params,
        "module": call.module,
        "foo": call.foo,
        "input": tuple(call.input) if call.input else (),
        "output": tuple(call.output) if call.output else (),
        "input_groups": (
            tuple(tuple(group) for group in call.input_groups)
            if call.input_groups
            else ()
        ),
        "output_groups": (
            tuple(tuple(group) for group in call.output_groups)
            if call.output_groups
            else ()
        ),
        "dress": call.dress,
        "strategy": call.strategy or "single",
    }
    if inline_target is not None:
        transform_kwargs["target"] = inline_target
    if call.use is None and call.keys is not None:
        transform_kwargs["keys"] = KeySelectionConfig.model_validate(
            call.keys.model_dump()
        )
    # When call.use references ingestion_model.transforms, defer strict
    # transform validation until finish_init can hydrate module/foo.
    if call.use is not None and call.module is None and call.foo is None:
        self.t = Transform(name=call.use)
        return
    self.t = Transform(**transform_kwargs)
fetch_important_items()
Source code in graflo/architecture/pipeline/runtime/actor/transform.py
def fetch_important_items(self) -> dict[str, Any]:
    items = self._fetch_items_from_dict(("transform",))
    items.update({"t.input": self.t.input, "t.output": self.t.output})
    return items
finish_init(init_ctx)
Source code in graflo/architecture/pipeline/runtime/actor/transform.py
def finish_init(self, init_ctx: ActorInitContext) -> None:
    self.transforms = init_ctx.transforms
    if self.call_use is None or self.t._foo is not None:
        self._refresh_missing_key_guard(init_ctx)
        return
    if self._call_config is None:
        self._refresh_missing_key_guard(init_ctx)
        return
    pt = self.transforms.get(self.call_use, None)
    if pt is None:
        if init_ctx.strict_references:
            raise ValueError(
                f"Transform '{self.call_use}' referenced by transform.call.use "
                "was not found in ingestion_model.transforms."
            )
        self._refresh_missing_key_guard(init_ctx)
        return
    call = self._call_config
    transform_kwargs = self._merge_call_with_proto(call, pt)
    self.t = Transform(**transform_kwargs)
    self._refresh_missing_key_guard(init_ctx)
from_config(config) classmethod
Source code in graflo/architecture/pipeline/runtime/actor/transform.py
@classmethod
def from_config(cls, config: TransformActorConfig) -> TransformActor:
    return cls(config)
init_transforms(init_ctx)
Source code in graflo/architecture/pipeline/runtime/actor/transform.py
def init_transforms(self, init_ctx: ActorInitContext) -> None:
    self.transforms = init_ctx.transforms
references_vertices()
Source code in graflo/architecture/pipeline/runtime/actor/transform.py
def references_vertices(self) -> set[str]:
    return set()

VertexActor

Bases: VertexProducingActor

Actor for processing vertex data.

Source code in graflo/architecture/pipeline/runtime/actor/vertex.py
class VertexActor(VertexProducingActor):
    """Actor for processing vertex data."""

    def __init__(self, config: VertexActorConfig):
        self.name = config.vertex
        self.from_doc: dict[str, str] | None = config.from_doc
        self.keep_fields: tuple[str, ...] | None = (
            tuple(config.keep_fields) if config.keep_fields else None
        )
        self.extraction_scope: Literal["full", "mapped_only"] = config.extraction_scope
        self.role: str | None = config.role
        self.lookup_only: bool = config.lookup_only
        self.vertex_config: VertexConfig
        self.allowed_vertex_names: set[VertexName] | None = None

    @classmethod
    def from_config(cls, config: VertexActorConfig) -> VertexActor:
        return cls(config)

    def fetch_important_items(self) -> dict[str, Any]:
        return self._fetch_items_from_dict(
            (
                "name",
                "from_doc",
                "keep_fields",
                "extraction_scope",
                "role",
                "lookup_only",
            )
        )

    def finish_init(self, init_ctx: ActorInitContext) -> None:
        self.vertex_config = init_ctx.vertex_config
        self.allowed_vertex_names = init_ctx.allowed_vertex_names
        if init_ctx.strict_references and self.from_doc:
            refuse_undeclared_mapping(self.vertex_config, self.name, self.from_doc)

    def _filter_and_aggregate_vertex_docs(
        self, docs: list[dict[str, Any]], doc: dict[str, Any]
    ) -> list[dict[str, Any]]:
        filters = self.vertex_config.filters(self.name)
        return [
            _doc for _doc in docs if all(cfilter.matches(_doc) for cfilter in filters)
        ]

    def _extract_vertex_doc_from_transformed_item(
        self,
        item: Any,
        vertex_keys: tuple[str, ...],
        index_keys: tuple[str, ...],
    ) -> dict[str, Any]:
        if isinstance(item, TransformPayload):
            doc: dict[str, Any] = {}
            consumed_named: set[str] = set()
            for k, v in item.named.items():
                if k in vertex_keys and v is not None:
                    doc[k] = v
                    consumed_named.add(k)
            for j, value in enumerate(item.positional):
                if j >= len(index_keys):
                    break
                doc[index_keys[j]] = value
            for key in consumed_named:
                item.named.pop(key, None)
            if item.positional:
                item.positional = ()
            return doc

        if isinstance(item, dict):
            doc = {}
            value_keys = sorted(
                (
                    k
                    for k in item
                    if k.startswith(ActorConstants.DRESSING_TRANSFORMED_VALUE_KEY)
                ),
                key=lambda x: int(x.rsplit("#", 1)[-1]),
            )
            for j, vkey in enumerate(value_keys):
                if j >= len(index_keys):
                    break
                doc[index_keys[j]] = item.pop(vkey)
            for vkey in vertex_keys:
                if vkey not in doc and vkey in item and item[vkey] is not None:
                    doc[vkey] = item.pop(vkey)
            return doc

        return {}

    def _process_transformed_items(
        self,
        ctx: ExtractionContext,
        lindex: LocationIndex,
        doc: dict[str, Any],
        vertex_keys: tuple[str, ...],
    ) -> list[dict[str, Any]]:
        index_keys = tuple(self.vertex_config.identity_fields(self.name))
        payloads = ctx.transform_buffer[lindex]
        extracted_docs = [
            self._extract_vertex_doc_from_transformed_item(
                item, vertex_keys, index_keys
            )
            for item in payloads
        ]
        ctx.transform_buffer[lindex] = [
            item
            for item in payloads
            if not (
                isinstance(item, TransformPayload)
                and not item.named
                and not item.positional
            )
            and not (isinstance(item, dict) and not item)
        ]
        return self._filter_and_aggregate_vertex_docs(extracted_docs, doc)

    def __call__(
        self, ctx: ExtractionContext, lindex: LocationIndex, *nargs: Any, **kwargs: Any
    ) -> ExtractionContext:
        doc: dict[str, Any] = kwargs.get("doc", {})
        buffer_items: list[Any] = list(ctx.transform_buffer.get(lindex, []))
        effective_doc = merge_observation_with_transform_buffer(doc, buffer_items)
        ctx.obs_buffer[lindex] = dict(effective_doc)

        # Early-exit for disallowed vertex types.
        # This must happen before any ctx.acc_vertex[...] access.
        if (
            self.allowed_vertex_names is not None
            and self.name not in self.allowed_vertex_names
        ):
            return ctx
        if (
            self.allowed_vertex_names is None
            and self.name not in self.vertex_config.vertex_set
        ):
            return ctx

        vertex_keys_list = self.vertex_config.property_names(self.name)
        vertex_keys: tuple[str, ...] = tuple(vertex_keys_list)

        # When a role is set the vertex is stored at a named sub-slot so that
        # multiple vertices of the same type in one observation (e.g. buyer/seller)
        # occupy distinct accumulator locations. Transforms are always read from
        # the bare observation lindex; only storage moves to the role slot.
        effective_lindex = lindex.extend((self.role, 0)) if self.role else lindex

        agg = []
        identity_fields = self.vertex_config.identity_fields(self.name)
        if self.from_doc:
            source_keys = set(self.from_doc.values())
            consumed_from_buffer = False
            for item in ctx.transform_buffer[lindex]:
                if isinstance(item, TransformPayload) and source_keys.issubset(
                    item.named
                ):
                    projected = {
                        v_f: item.named[d_f] for v_f, d_f in self.from_doc.items()
                    }
                    if any(v is not None for v in projected.values()):
                        agg.extend(explode_identity_lists(projected, identity_fields))
                    for k in source_keys:
                        item.named.pop(k, None)
                    consumed_from_buffer = True
            ctx.transform_buffer[lindex] = [
                item
                for item in ctx.transform_buffer[lindex]
                if not (
                    isinstance(item, TransformPayload)
                    and not item.named
                    and not item.positional
                )
                and not (isinstance(item, dict) and not item)
            ]
            if not consumed_from_buffer:
                projected = {
                    v_f: effective_doc.get(d_f) for v_f, d_f in self.from_doc.items()
                }
                if any(v is not None for v in projected.values()):
                    agg.extend(explode_identity_lists(projected, identity_fields))
            buffer_vertex_keys = tuple(k for k in vertex_keys if k not in self.from_doc)
        else:
            buffer_vertex_keys = vertex_keys

        agg.extend(
            self._process_transformed_items(
                ctx, lindex, effective_doc, buffer_vertex_keys
            )
        )

        if self.extraction_scope == "full":
            remaining_keys = set(vertex_keys) - set().union(*[d.keys() for d in agg])
            # When keep_fields is set, restrict passthrough to only those declared fields.
            if self.keep_fields is not None:
                remaining_keys = remaining_keys & set(self.keep_fields)
            passthrough_doc = {
                k: effective_doc.get(k) for k in remaining_keys if k in effective_doc
            }
            if passthrough_doc:
                agg.append(passthrough_doc)

        if self.name in self.vertex_config.hash_identity_vertices:
            # A digest vertex's key is synthetic: a value the record carries
            # under that name (a source `id` column, a mapping, a positional
            # transform output) would key the document instead of the digest.
            for vertex_doc in agg:
                for field in identity_fields:
                    vertex_doc.pop(field, None)

        merged = fuse_doc_basis(agg, index_keys=tuple(identity_fields))

        for m in merged:
            vertex_rep = VertexRep(vertex=m, lookup_only=self.lookup_only)
            ctx.acc_vertex[self.name][effective_lindex].append(vertex_rep)
            ctx.record_vertex_observation(
                vertex_name=self.name,
                location=effective_lindex,
                vertex=vertex_rep.vertex,
                ctx={},
            )
        return ctx

    def references_vertices(self) -> set[VertexName]:
        return {self.name}

Attributes

allowed_vertex_names = None instance-attribute
extraction_scope = config.extraction_scope instance-attribute
from_doc = config.from_doc instance-attribute
keep_fields = tuple(config.keep_fields) if config.keep_fields else None instance-attribute
lookup_only = config.lookup_only instance-attribute
name = config.vertex instance-attribute
role = config.role instance-attribute
vertex_config instance-attribute

Methods:

__call__(ctx, lindex, *nargs, **kwargs)
Source code in graflo/architecture/pipeline/runtime/actor/vertex.py
def __call__(
    self, ctx: ExtractionContext, lindex: LocationIndex, *nargs: Any, **kwargs: Any
) -> ExtractionContext:
    doc: dict[str, Any] = kwargs.get("doc", {})
    buffer_items: list[Any] = list(ctx.transform_buffer.get(lindex, []))
    effective_doc = merge_observation_with_transform_buffer(doc, buffer_items)
    ctx.obs_buffer[lindex] = dict(effective_doc)

    # Early-exit for disallowed vertex types.
    # This must happen before any ctx.acc_vertex[...] access.
    if (
        self.allowed_vertex_names is not None
        and self.name not in self.allowed_vertex_names
    ):
        return ctx
    if (
        self.allowed_vertex_names is None
        and self.name not in self.vertex_config.vertex_set
    ):
        return ctx

    vertex_keys_list = self.vertex_config.property_names(self.name)
    vertex_keys: tuple[str, ...] = tuple(vertex_keys_list)

    # When a role is set the vertex is stored at a named sub-slot so that
    # multiple vertices of the same type in one observation (e.g. buyer/seller)
    # occupy distinct accumulator locations. Transforms are always read from
    # the bare observation lindex; only storage moves to the role slot.
    effective_lindex = lindex.extend((self.role, 0)) if self.role else lindex

    agg = []
    identity_fields = self.vertex_config.identity_fields(self.name)
    if self.from_doc:
        source_keys = set(self.from_doc.values())
        consumed_from_buffer = False
        for item in ctx.transform_buffer[lindex]:
            if isinstance(item, TransformPayload) and source_keys.issubset(
                item.named
            ):
                projected = {
                    v_f: item.named[d_f] for v_f, d_f in self.from_doc.items()
                }
                if any(v is not None for v in projected.values()):
                    agg.extend(explode_identity_lists(projected, identity_fields))
                for k in source_keys:
                    item.named.pop(k, None)
                consumed_from_buffer = True
        ctx.transform_buffer[lindex] = [
            item
            for item in ctx.transform_buffer[lindex]
            if not (
                isinstance(item, TransformPayload)
                and not item.named
                and not item.positional
            )
            and not (isinstance(item, dict) and not item)
        ]
        if not consumed_from_buffer:
            projected = {
                v_f: effective_doc.get(d_f) for v_f, d_f in self.from_doc.items()
            }
            if any(v is not None for v in projected.values()):
                agg.extend(explode_identity_lists(projected, identity_fields))
        buffer_vertex_keys = tuple(k for k in vertex_keys if k not in self.from_doc)
    else:
        buffer_vertex_keys = vertex_keys

    agg.extend(
        self._process_transformed_items(
            ctx, lindex, effective_doc, buffer_vertex_keys
        )
    )

    if self.extraction_scope == "full":
        remaining_keys = set(vertex_keys) - set().union(*[d.keys() for d in agg])
        # When keep_fields is set, restrict passthrough to only those declared fields.
        if self.keep_fields is not None:
            remaining_keys = remaining_keys & set(self.keep_fields)
        passthrough_doc = {
            k: effective_doc.get(k) for k in remaining_keys if k in effective_doc
        }
        if passthrough_doc:
            agg.append(passthrough_doc)

    if self.name in self.vertex_config.hash_identity_vertices:
        # A digest vertex's key is synthetic: a value the record carries
        # under that name (a source `id` column, a mapping, a positional
        # transform output) would key the document instead of the digest.
        for vertex_doc in agg:
            for field in identity_fields:
                vertex_doc.pop(field, None)

    merged = fuse_doc_basis(agg, index_keys=tuple(identity_fields))

    for m in merged:
        vertex_rep = VertexRep(vertex=m, lookup_only=self.lookup_only)
        ctx.acc_vertex[self.name][effective_lindex].append(vertex_rep)
        ctx.record_vertex_observation(
            vertex_name=self.name,
            location=effective_lindex,
            vertex=vertex_rep.vertex,
            ctx={},
        )
    return ctx
__init__(config)
Source code in graflo/architecture/pipeline/runtime/actor/vertex.py
def __init__(self, config: VertexActorConfig):
    self.name = config.vertex
    self.from_doc: dict[str, str] | None = config.from_doc
    self.keep_fields: tuple[str, ...] | None = (
        tuple(config.keep_fields) if config.keep_fields else None
    )
    self.extraction_scope: Literal["full", "mapped_only"] = config.extraction_scope
    self.role: str | None = config.role
    self.lookup_only: bool = config.lookup_only
    self.vertex_config: VertexConfig
    self.allowed_vertex_names: set[VertexName] | None = None
fetch_important_items()
Source code in graflo/architecture/pipeline/runtime/actor/vertex.py
def fetch_important_items(self) -> dict[str, Any]:
    return self._fetch_items_from_dict(
        (
            "name",
            "from_doc",
            "keep_fields",
            "extraction_scope",
            "role",
            "lookup_only",
        )
    )
finish_init(init_ctx)
Source code in graflo/architecture/pipeline/runtime/actor/vertex.py
def finish_init(self, init_ctx: ActorInitContext) -> None:
    self.vertex_config = init_ctx.vertex_config
    self.allowed_vertex_names = init_ctx.allowed_vertex_names
    if init_ctx.strict_references and self.from_doc:
        refuse_undeclared_mapping(self.vertex_config, self.name, self.from_doc)
from_config(config) classmethod
Source code in graflo/architecture/pipeline/runtime/actor/vertex.py
@classmethod
def from_config(cls, config: VertexActorConfig) -> VertexActor:
    return cls(config)
references_vertices()
Source code in graflo/architecture/pipeline/runtime/actor/vertex.py
def references_vertices(self) -> set[VertexName]:
    return {self.name}

VertexRouterActor

Bases: VertexProducingActor

Routes documents to the correct VertexActor based on a type field.

The merged observation (document + same-location transform buffer) is passed through to the selected :class:VertexActor unchanged. Projection uses the same from / vertex_from_map contract as a standalone vertex step.

Vertices are accumulated at lindex.extend((role, 0)). role is normalized by config validation (defaults to :attr:type_field when omitted), so runtime slot addressing uses a single internal key. A downstream dynamic EdgeActor references this slot via source_role / target_role (or source_type_field / target_type_field) using the same segment name.

lookup_only (every routed class, or the listed ones) is handed to the vertex actor of each class it covers, so those rows locate edge endpoints and are never written. vertex_types bounds the classes it produces: a value resolving to any other class is skipped.

Source code in graflo/architecture/pipeline/runtime/actor/vertex_router.py
class VertexRouterActor(VertexProducingActor):
    """Routes documents to the correct VertexActor based on a type field.

    The merged observation (document + same-location transform buffer) is passed
    through to the selected :class:`VertexActor` unchanged. Projection uses the same
    ``from`` / ``vertex_from_map`` contract as a standalone vertex step.

    Vertices are accumulated at ``lindex.extend((role, 0))``. ``role`` is normalized
    by config validation (defaults to :attr:`type_field` when omitted), so runtime slot
    addressing uses a single internal key. A downstream dynamic ``EdgeActor`` references
    this slot via ``source_role`` / ``target_role`` (or ``source_type_field`` /
    ``target_type_field``) using the same segment name.

    ``lookup_only`` (every routed class, or the listed ones) is handed to the
    vertex actor of each class it covers, so those rows locate edge endpoints
    and are never written. ``vertex_types`` bounds the classes it produces: a
    value resolving to any other class is skipped.
    """

    def __init__(self, config: VertexRouterActorConfig):
        self.config = config
        self.type_field = config.type_field
        # Config normalization guarantees role is always present.
        self.role: str = config.role or config.type_field
        self.from_doc: dict[str, str] | None = config.from_doc
        self.keep_fields: tuple[str, ...] | None = (
            tuple(config.keep_fields) if config.keep_fields else None
        )
        self.extraction_scope: Literal["full", "mapped_only"] = config.extraction_scope
        self.type_map: dict[str, str] = config.type_map or {}
        self.vertex_from_map: dict[str, dict[str, str]] = config.vertex_from_map or {}
        self._vertex_actors: dict[str, ActorWrapper] = {}
        self._init_ctx: ActorInitContext | None = None
        self.vertex_config: VertexConfig = VertexConfig(vertices=[])

    @classmethod
    def from_config(cls, config: VertexRouterActorConfig) -> VertexRouterActor:
        return cls(config)

    def fetch_important_items(self) -> dict[str, Any]:
        items: dict[str, Any] = {"type_field": self.type_field, "role": self.role}
        if self.from_doc:
            items["from_doc"] = self.from_doc
        if self.keep_fields:
            items["keep_fields"] = list(self.keep_fields)
        items["extraction_scope"] = self.extraction_scope
        if self.type_map:
            items["type_map"] = self.type_map
        if self.vertex_from_map:
            items["vertex_from_map"] = self.vertex_from_map
        if self.config.type_map_only:
            items["type_map_only"] = True
        if self.config.vertex_types is not None:
            items["vertex_types"] = self.config.vertex_types
        if self.config.lookup_only:
            items["lookup_only"] = self.config.lookup_only
        items["routed_types"] = sorted(self._vertex_actors.keys())
        return items

    def finish_init(self, init_ctx: ActorInitContext) -> None:
        self.vertex_config = init_ctx.vertex_config
        self._init_ctx = init_ctx
        self._vertex_actors.clear()
        for field, named in (
            ("lookup_only", self.config.lookup_only),
            ("vertex_types", self.config.vertex_types),
        ):
            if isinstance(named, list):
                unknown = sorted(set(named) - self.vertex_config.vertex_set)
                if unknown:
                    raise ValueError(
                        f"vertex_router on {self.type_field!r}: {field} names "
                        f"{unknown}, which the schema does not declare"
                    )
        if init_ctx.strict_references:
            # Checked for the classes the router names; one an unbounded router
            # routes to by pass-through is known only per record, and gets the
            # declared part of the map (see `_get_or_create_wrapper`).
            for vertex_type, from_doc in self.vertex_from_map.items():
                refuse_undeclared_mapping(self.vertex_config, vertex_type, from_doc)
            if self.from_doc:
                for vertex_type in sorted(self._named_classes()):
                    if vertex_type not in self.vertex_from_map:
                        refuse_undeclared_mapping(
                            self.vertex_config, vertex_type, self.from_doc
                        )

    def _named_classes(self) -> set[str]:
        """The classes the router can produce that its config names.

        The table's classes, bounded by ``vertex_types``; an open bounded router
        reaches every listed class by pass-through as well.
        """
        table = set(self.type_map.values())
        bound = self.config.vertex_types
        if bound is None:
            return table
        if self.config.type_map_only:
            return table & set(bound)
        return set(bound)

    def _get_or_create_wrapper(self, vertex_type: str) -> ActorWrapper | None:
        from .wrapper import ActorWrapper

        if vertex_type not in self.vertex_config.vertex_set:
            return None
        # Fast path without the lock: dict reads are atomic, and a wrapper is only
        # published after finish_init completes, so it is never seen half-built.
        wrapper = self._vertex_actors.get(vertex_type)
        if wrapper is not None:
            return wrapper
        if self._init_ctx is None:
            raise RuntimeError(
                "VertexRouterActor._get_or_create_wrapper called before finish_init"
            )

        with _VERTEX_ACTOR_CREATE_LOCK:
            wrapper = self._vertex_actors.get(vertex_type)
            if wrapper is not None:
                return wrapper
            if vertex_type in self.vertex_from_map:
                per_type_from = self.vertex_from_map[vertex_type]
            else:
                per_type_from = self.from_doc
            if self._init_ctx.strict_references and per_type_from:
                per_type_from = self._declared_part(vertex_type, per_type_from)
            config = VertexActorConfig.model_validate(
                {
                    "type": "vertex",
                    "vertex": vertex_type,
                    "from": per_type_from,
                    "keep_fields": list(self.keep_fields) if self.keep_fields else None,
                    "extraction_scope": self.extraction_scope,
                    "lookup_only": self.config.looks_up(vertex_type),
                }
            )
            wrapper = ActorWrapper.from_config(config)
            wrapper.finish_init(self._init_ctx)
            self._vertex_actors[vertex_type] = wrapper
        logger.debug(
            "VertexRouterActor: lazily registered VertexActor(%s) for type_field=%s role=%s",
            vertex_type,
            self.type_field,
            self.role,
        )
        return wrapper

    def _declared_part(
        self, vertex_type: str, from_doc: dict[str, str]
    ) -> dict[str, str]:
        """*from_doc* without the targets *vertex_type* does not declare, said once."""
        unknown = undeclared_mapping_targets(self.vertex_config, vertex_type, from_doc)
        if not unknown:
            return from_doc
        logger.warning(
            "vertex_router on %r: %s routed by pass-through does not declare %s; "
            "those `from` targets are not mapped",
            self.type_field,
            vertex_type,
            unknown,
        )
        return {k: v for k, v in from_doc.items() if k not in unknown}

    def count(self) -> int:
        return 1 + sum(w.count() for w in self._vertex_actors.values())

    def references_vertices(self) -> set[VertexName]:
        return set(self._vertex_actors.keys())

    def __call__(
        self, ctx: ExtractionContext, lindex: LocationIndex, *nargs: Any, **kwargs: Any
    ) -> ExtractionContext:
        raw_observation = kwargs.get("doc", {})
        if not isinstance(raw_observation, dict):
            logger.debug(
                "VertexRouterActor: expected dict observation slice, got %s, skipping",
                type(raw_observation).__name__,
            )
            return ctx
        buffer_items: list[Any] = list(ctx.transform_buffer.get(lindex, []))
        doc = merge_observation_with_transform_buffer(raw_observation, buffer_items)
        ctx.obs_buffer[lindex] = dict(doc)
        raw_vtype = doc.get(self.type_field)
        if raw_vtype is None:
            logger.debug(
                "VertexRouterActor: type_field '%s' not in doc, skipping",
                self.type_field,
            )
            return ctx
        if self.config.type_map_only and raw_vtype not in self.type_map:
            logger.debug(
                "VertexRouterActor: value %r of '%s' is not in the closed "
                "type_map, skipping",
                raw_vtype,
                self.type_field,
            )
            return ctx
        vtype = self.type_map.get(raw_vtype, raw_vtype)
        if not self.config.admits(vtype):
            logger.debug(
                "VertexRouterActor: vertex type %r (from field '%s') is outside "
                "vertex_types, skipping",
                vtype,
                self.type_field,
            )
            return ctx

        wrapper = self._get_or_create_wrapper(vtype)
        if wrapper is None:
            logger.debug(
                "VertexRouterActor: vertex type '%s' (from field '%s') "
                "not in VertexConfig, skipping",
                vtype,
                self.type_field,
            )
            return ctx

        effective_lindex = lindex.extend((self.role, 0))
        return wrapper(ctx, effective_lindex, doc=doc)

Attributes

config = config instance-attribute
extraction_scope = config.extraction_scope instance-attribute
from_doc = config.from_doc instance-attribute
keep_fields = tuple(config.keep_fields) if config.keep_fields else None instance-attribute
role = config.role or config.type_field instance-attribute
type_field = config.type_field instance-attribute
type_map = config.type_map or {} instance-attribute
vertex_config = VertexConfig(vertices=[]) instance-attribute
vertex_from_map = config.vertex_from_map or {} instance-attribute

Methods:

__call__(ctx, lindex, *nargs, **kwargs)
Source code in graflo/architecture/pipeline/runtime/actor/vertex_router.py
def __call__(
    self, ctx: ExtractionContext, lindex: LocationIndex, *nargs: Any, **kwargs: Any
) -> ExtractionContext:
    raw_observation = kwargs.get("doc", {})
    if not isinstance(raw_observation, dict):
        logger.debug(
            "VertexRouterActor: expected dict observation slice, got %s, skipping",
            type(raw_observation).__name__,
        )
        return ctx
    buffer_items: list[Any] = list(ctx.transform_buffer.get(lindex, []))
    doc = merge_observation_with_transform_buffer(raw_observation, buffer_items)
    ctx.obs_buffer[lindex] = dict(doc)
    raw_vtype = doc.get(self.type_field)
    if raw_vtype is None:
        logger.debug(
            "VertexRouterActor: type_field '%s' not in doc, skipping",
            self.type_field,
        )
        return ctx
    if self.config.type_map_only and raw_vtype not in self.type_map:
        logger.debug(
            "VertexRouterActor: value %r of '%s' is not in the closed "
            "type_map, skipping",
            raw_vtype,
            self.type_field,
        )
        return ctx
    vtype = self.type_map.get(raw_vtype, raw_vtype)
    if not self.config.admits(vtype):
        logger.debug(
            "VertexRouterActor: vertex type %r (from field '%s') is outside "
            "vertex_types, skipping",
            vtype,
            self.type_field,
        )
        return ctx

    wrapper = self._get_or_create_wrapper(vtype)
    if wrapper is None:
        logger.debug(
            "VertexRouterActor: vertex type '%s' (from field '%s') "
            "not in VertexConfig, skipping",
            vtype,
            self.type_field,
        )
        return ctx

    effective_lindex = lindex.extend((self.role, 0))
    return wrapper(ctx, effective_lindex, doc=doc)
__init__(config)
Source code in graflo/architecture/pipeline/runtime/actor/vertex_router.py
def __init__(self, config: VertexRouterActorConfig):
    self.config = config
    self.type_field = config.type_field
    # Config normalization guarantees role is always present.
    self.role: str = config.role or config.type_field
    self.from_doc: dict[str, str] | None = config.from_doc
    self.keep_fields: tuple[str, ...] | None = (
        tuple(config.keep_fields) if config.keep_fields else None
    )
    self.extraction_scope: Literal["full", "mapped_only"] = config.extraction_scope
    self.type_map: dict[str, str] = config.type_map or {}
    self.vertex_from_map: dict[str, dict[str, str]] = config.vertex_from_map or {}
    self._vertex_actors: dict[str, ActorWrapper] = {}
    self._init_ctx: ActorInitContext | None = None
    self.vertex_config: VertexConfig = VertexConfig(vertices=[])
count()
Source code in graflo/architecture/pipeline/runtime/actor/vertex_router.py
def count(self) -> int:
    return 1 + sum(w.count() for w in self._vertex_actors.values())
fetch_important_items()
Source code in graflo/architecture/pipeline/runtime/actor/vertex_router.py
def fetch_important_items(self) -> dict[str, Any]:
    items: dict[str, Any] = {"type_field": self.type_field, "role": self.role}
    if self.from_doc:
        items["from_doc"] = self.from_doc
    if self.keep_fields:
        items["keep_fields"] = list(self.keep_fields)
    items["extraction_scope"] = self.extraction_scope
    if self.type_map:
        items["type_map"] = self.type_map
    if self.vertex_from_map:
        items["vertex_from_map"] = self.vertex_from_map
    if self.config.type_map_only:
        items["type_map_only"] = True
    if self.config.vertex_types is not None:
        items["vertex_types"] = self.config.vertex_types
    if self.config.lookup_only:
        items["lookup_only"] = self.config.lookup_only
    items["routed_types"] = sorted(self._vertex_actors.keys())
    return items
finish_init(init_ctx)
Source code in graflo/architecture/pipeline/runtime/actor/vertex_router.py
def finish_init(self, init_ctx: ActorInitContext) -> None:
    self.vertex_config = init_ctx.vertex_config
    self._init_ctx = init_ctx
    self._vertex_actors.clear()
    for field, named in (
        ("lookup_only", self.config.lookup_only),
        ("vertex_types", self.config.vertex_types),
    ):
        if isinstance(named, list):
            unknown = sorted(set(named) - self.vertex_config.vertex_set)
            if unknown:
                raise ValueError(
                    f"vertex_router on {self.type_field!r}: {field} names "
                    f"{unknown}, which the schema does not declare"
                )
    if init_ctx.strict_references:
        # Checked for the classes the router names; one an unbounded router
        # routes to by pass-through is known only per record, and gets the
        # declared part of the map (see `_get_or_create_wrapper`).
        for vertex_type, from_doc in self.vertex_from_map.items():
            refuse_undeclared_mapping(self.vertex_config, vertex_type, from_doc)
        if self.from_doc:
            for vertex_type in sorted(self._named_classes()):
                if vertex_type not in self.vertex_from_map:
                    refuse_undeclared_mapping(
                        self.vertex_config, vertex_type, self.from_doc
                    )
from_config(config) classmethod
Source code in graflo/architecture/pipeline/runtime/actor/vertex_router.py
@classmethod
def from_config(cls, config: VertexRouterActorConfig) -> VertexRouterActor:
    return cls(config)
references_vertices()
Source code in graflo/architecture/pipeline/runtime/actor/vertex_router.py
def references_vertices(self) -> set[VertexName]:
    return set(self._vertex_actors.keys())