Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
136 changes: 136 additions & 0 deletions src/flowx/discovery_lineage.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,136 @@
"""Source-neutral lineage derivation over the shared discovery AST.

The parallel of :mod:`flowx.lineage`, which derives a
:class:`~flowx.models.ir.Lineage` block over the Databricks IR. This module
derives the same block over the shared discovery AST
(:mod:`flowx.models.discovery`) instead, so the discover phase can attach lineage
to a :class:`~flowx.models.discovery.SourceGraph` before any IR translation
exists.

It reuses :mod:`flowx.lineage`'s primitive cores -- :func:`control_edges_from_calls`
and :func:`data_edges_from_endpoints` -- so the two-tier match, the self-edge drop,
and the dedup live in exactly one place and both phases behave identically. It
imports nothing from ``sources/*``: the walk is over the neutral
:class:`~flowx.models.discovery.ContainerNode` branch shape, so an ADF or an
Airflow graph derives lineage through this one code path once its nodes carry
``data_reads`` / ``data_writes`` and (for control edges) the
:data:`INVOKES_WORKFLOW_PROPERTY` marker.
"""

from __future__ import annotations

import dataclasses
from collections.abc import Iterator

from flowx.lineage import control_edges_from_calls, data_edges_from_endpoints
from flowx.models.discovery import ContainerNode, SourceGraph, SourceNode
from flowx.models.ir import ControlEdge, DataEdge, Lineage

# Neutral node-property key under which a mapper records the workflow a node
# invokes (an ADF ``ExecutePipeline`` callee, an Airflow triggered job). Kept in
# the free-form ``properties`` seam because the invocation target is a per-source
# detail with no shared typed field; :func:`build_graph_control_edges` reads it
# here so the control-edge derivation stays source-agnostic.
INVOKES_WORKFLOW_PROPERTY = "invokes_workflow"
# Companion key: whether the caller waits for the invoked workflow to complete
# (``True`` / ``False``), or absent when the source has no such notion.
INVOKES_WAIT_PROPERTY = "invokes_wait"


def walk_nodes(nodes: list[SourceNode]) -> Iterator[SourceNode]:
"""Yield every node depth-first, descending into every container branch.

Recurses through :class:`ContainerNode` branches in their insertion order, so
a Switch's cases *and* its ``default`` branch, a ForEach / Until ``body``, and
both sides of an IfCondition are all reached -- a data asset or an invocation
buried inside a Switch case is still found.

Args:
nodes: Top-level (or already-nested) node list to walk.

Yields:
Each node, container nodes included, in depth-first order.
"""
for node in nodes:
yield node
if isinstance(node, ContainerNode):
for children in node.branches.values():
yield from walk_nodes(children)


def build_graph_control_edges(graph: SourceGraph) -> list[ControlEdge]:
"""Derive cross-workflow invocation edges for a discovery graph.

One edge per node that carries the :data:`INVOKES_WORKFLOW_PROPERTY` marker,
found anywhere in the graph (fan-out inside ForEach / If / Switch preserved).
Delegates the self-edge drop, dedup, and unresolved-callee recording to the
shared :func:`~flowx.lineage.control_edges_from_calls`.

Args:
graph: The source graph to derive control edges for.

Returns:
Deduplicated control edges, in first-seen order.
"""

def _calls() -> Iterator[tuple[str, bool | None, str]]:
for node in walk_nodes(graph.tasks):
if INVOKES_WORKFLOW_PROPERTY not in node.properties:
continue
target = node.properties.get(INVOKES_WORKFLOW_PROPERTY) or ""
wait = node.properties.get(INVOKES_WAIT_PROPERTY)
yield str(target), wait, node.task_key

return control_edges_from_calls(graph.name, _calls())


def build_graph_data_edges(graph: SourceGraph) -> list[DataEdge]:
"""Derive proven producer -> consumer data hand-offs for a discovery graph.

A producer is any node with a ``data_writes`` asset; a consumer any node with a
``data_reads`` asset, gathered across the whole graph (every container branch
included). Delegates the two-tier match, self-edge drop, and dedup to the shared
:func:`~flowx.lineage.data_edges_from_endpoints`.

Args:
graph: The source graph to derive data edges for.

Returns:
Deduplicated data edges, in first-seen order.
"""
nodes = list(walk_nodes(graph.tasks))
producers = [(node.task_key, asset) for node in nodes for asset in node.data_writes]
consumers = [(node.task_key, asset) for node in nodes for asset in node.data_reads]
return data_edges_from_endpoints(producers, consumers)


def build_graph_lineage(graph: SourceGraph) -> Lineage:
"""Compose the source-neutral lineage block for a discovery graph.

Motif annotations are a convert-time IR concern (motifs are detected during
translation, not discovery), so the discovery lineage block leaves them empty.

Args:
graph: The source graph to derive lineage for.

Returns:
A :class:`Lineage` with control edges and data edges (motifs empty).
"""
return Lineage(
control_edges=build_graph_control_edges(graph),
data_edges=build_graph_data_edges(graph),
)


def with_graph_lineage(graph: SourceGraph) -> SourceGraph:
"""Return a *new* graph carrying its derived lineage, leaving the input untouched.

Mirrors :func:`flowx.lineage.with_lineage` for the discovery AST.

Args:
graph: The source graph to copy.

Returns:
A shallow copy of *graph* with :attr:`SourceGraph.lineage` populated.
"""
return dataclasses.replace(graph, lineage=build_graph_lineage(graph))
53 changes: 51 additions & 2 deletions src/flowx/discovery_serde.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,10 @@

The DataAsset (de)serialisers are reused from ``ir_serde`` (``data_asset_to_dict``
/ ``data_asset_from_dict``) so the physical-asset shape has a single definition
shared by the lineage substrate and the discovery AST.
shared by the lineage substrate and the discovery AST. A graph's derived
:class:`~flowx.models.ir.Lineage` block is serialised through ``ir_serde``'s
``lineage_to_dict`` for the same reason; its inverse (:func:`_lineage_from_dict`)
lives here because ``ir_serde`` ships only the forward direction.

Every node dict carries a ``node_type`` discriminator (the dataclass name) so a
:class:`~flowx.models.discovery.ContainerNode` or
Expand All @@ -19,7 +22,7 @@

from typing import Any

from flowx.ir_serde import data_asset_from_dict, data_asset_to_dict
from flowx.ir_serde import data_asset_from_dict, data_asset_to_dict, lineage_to_dict
from flowx.models.discovery import (
ContainerNode,
GapNode,
Expand All @@ -30,6 +33,7 @@
SourceGraph,
SourceNode,
)
from flowx.models.ir import ControlEdge, DataEdge, Lineage, MotifAnnotation


def source_graph_to_dict(graph: SourceGraph) -> dict[str, Any]:
Expand All @@ -46,6 +50,8 @@ def source_graph_to_dict(graph: SourceGraph) -> dict[str, Any]:
result["description"] = graph.description
if graph.schedule is not None:
result["schedule"] = _schedule_to_dict(graph.schedule)
if graph.lineage is not None:
result["lineage"] = lineage_to_dict(graph.lineage)
if graph.properties:
result["properties"] = dict(graph.properties)
if graph.extensions:
Expand All @@ -58,6 +64,7 @@ def source_graph_to_dict(graph: SourceGraph) -> dict[str, Any]:
def source_graph_from_dict(raw: dict[str, Any]) -> SourceGraph:
"""Rehydrate a :class:`SourceGraph` from the dict :func:`source_graph_to_dict` emits."""
schedule = raw.get("schedule")
lineage = raw.get("lineage")
return SourceGraph(
name=raw.get("name", ""),
source=raw.get("source", ""),
Expand All @@ -67,12 +74,54 @@ def source_graph_from_dict(raw: dict[str, Any]) -> SourceGraph:
schedule=_schedule_from_dict(schedule) if schedule else None,
tags=list(raw.get("tags") or []),
tasks=[_node_from_dict(node) for node in raw.get("tasks") or []],
lineage=_lineage_from_dict(lineage) if lineage else None,
properties=dict(raw.get("properties") or {}),
extensions=dict(raw.get("extensions") or {}),
raw=raw.get("raw"),
)


def _lineage_from_dict(raw: dict[str, Any]) -> Lineage:
"""Rehydrate a :class:`Lineage` block from the dict ``ir_serde.lineage_to_dict`` emits.

The inverse of that forward serialiser (which ``ir_serde`` does not itself
ship), so a discovery graph's lineage round-trips through this module.
"""
return Lineage(
control_edges=[
ControlEdge(
source_workflow=edge.get("source_workflow", ""),
target_workflow=edge.get("target_workflow", ""),
via_task_key=edge.get("via_task_key", ""),
wait_for_completion=edge.get("wait_for_completion"),
resolved=bool(edge.get("resolved", True)),
)
for edge in raw.get("control_edges") or []
],
data_edges=[
DataEdge(
source_task_key=edge.get("source_task_key", ""),
target_task_key=edge.get("target_task_key", ""),
match_kind=edge.get("match_kind", ""),
match_key=edge.get("match_key", ""),
identity=edge.get("identity"),
asset_type=edge.get("asset_type"),
)
for edge in raw.get("data_edges") or []
],
motifs=[
MotifAnnotation(
motif_id=motif.get("motif_id", ""),
member_task_keys=list(motif.get("member_task_keys") or []),
display_name=motif.get("display_name"),
databricks_replacement=motif.get("databricks_replacement"),
notes=list(motif.get("notes") or []),
)
for motif in raw.get("motifs") or []
],
)


def _parameter_to_dict(spec: ParameterSpec) -> dict[str, Any]:
result: dict[str, Any] = {}
if spec.type is not None:
Expand Down
Loading