Skip to content
Closed
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
66 changes: 65 additions & 1 deletion src/flowx/bundler/dab_writer.py
Original file line number Diff line number Diff line change
Expand Up @@ -30,16 +30,21 @@
from flowx.models.ir import (
Activity,
AppendVariableActivity,
ControlEdge,
CopyActivity,
DataAsset,
DataEdge,
DbtFactoryActivity,
DeleteActivity,
Dependency,
ExecutePipelineActivity,
FilterActivity,
ForEachActivity,
IfConditionActivity,
Lineage,
LookupActivity,
MotifActivity,
MotifAnnotation,
NotebookActivity,
Pipeline,
PlaceholderActivity,
Expand Down Expand Up @@ -2037,6 +2042,7 @@ def pipeline_dict_to_ir(pipeline_dict: dict[str, Any]) -> tuple[Pipeline, list[d
reconciliation_status=pipeline_dict.get("reconciliation_status"),
migration_status=pipeline_dict.get("migration_status", "included"),
audit=dict(pipeline_dict.get("audit") or {}),
lineage=_reconstruct_lineage(pipeline_dict.get("lineage")),
)
return pipeline, parameters

Expand Down Expand Up @@ -2258,9 +2264,10 @@ def _reconstruct_ir(task_ir: dict[str, Any]) -> Activity:
bridge_required_parameters=dict(task_ir.get("bridge_required_parameters") or {}),
)
if task_type == "MotifActivity":
# motif_id arrives via ``base`` (Activity now owns the field); passing it again here
# would raise "multiple values for keyword argument 'motif_id'".
return MotifActivity(
**base,
motif_id=task_ir.get("motif_id", "unknown"),
display_name=task_ir.get("display_name", base["name"]),
databricks_replacement=task_ir.get("databricks_replacement", "notebook"),
matched_activity_names=list(task_ir.get("matched_activity_names", [])),
Expand Down Expand Up @@ -2311,9 +2318,66 @@ def _common_activity_kwargs(task_ir: dict[str, Any]) -> dict[str, Any]:
"required_parameters": dict(task_ir.get("required_parameters") or {}),
"compute_mode": task_ir.get("compute_mode"),
"notifications": task_ir.get("notifications"),
"motif_id": task_ir.get("motif_id"),
"data_reads": _reconstruct_data_assets(task_ir.get("data_reads")),
"data_writes": _reconstruct_data_assets(task_ir.get("data_writes")),
}


def _reconstruct_data_assets(raw: list[dict[str, Any]] | None) -> list[DataAsset]:
"""Rehydrates serialised DataAsset dicts into typed :class:`DataAsset` nodes."""
if not raw:
return []
return [
DataAsset(
signature=asset.get("signature", ""),
identity=asset.get("identity"),
asset_type=asset.get("asset_type"),
properties=dict(asset.get("properties") or {}),
)
for asset in raw
]


def _reconstruct_lineage(raw: dict[str, Any] | None) -> Lineage | None:
"""Rehydrates a serialised lineage block into a typed :class:`Lineage`, or ``None``."""
if not raw:
return None
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 _reconstruct_dependencies(raw: list[dict[str, Any]] | None) -> list[Dependency] | None:
if not raw:
return None
Expand Down
10 changes: 9 additions & 1 deletion src/flowx/ir_serde.py
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,8 @@ def pipeline_to_dict(pipeline: Pipeline) -> dict[str, Any]:
}
if pipeline.translation_configuration is not None:
result["translation_configuration"] = configuration_to_dict(pipeline.translation_configuration)
if pipeline.lineage is not None:
result["lineage"] = lineage_to_dict(pipeline.lineage)
return result


Expand Down Expand Up @@ -216,6 +218,12 @@ def activity_to_dict(task: Activity) -> dict[str, Any]:
task_dict["libraries"] = task.libraries
if task.parameter_approximations:
task_dict["parameter_approximations"] = task.parameter_approximations
if task.motif_id:
task_dict["motif_id"] = task.motif_id
if task.data_reads:
task_dict["data_reads"] = [data_asset_to_dict(asset) for asset in task.data_reads]
if task.data_writes:
task_dict["data_writes"] = [data_asset_to_dict(asset) for asset in task.data_writes]

extra = activity_extra_fields(task)
task_dict.update(extra)
Expand Down Expand Up @@ -410,7 +418,7 @@ def activity_extra_fields(activity: Activity) -> dict[str, Any]:
if activity.job_parameters:
extra["job_parameters"] = activity.job_parameters
case MotifActivity():
extra["motif_id"] = activity.motif_id
# motif_id now lives on the Activity base and is serialised by activity_to_dict.
extra["display_name"] = activity.display_name
extra["databricks_replacement"] = activity.databricks_replacement
extra["matched_activity_names"] = activity.matched_activity_names
Expand Down
217 changes: 193 additions & 24 deletions src/flowx/lineage.py
Original file line number Diff line number Diff line change
@@ -1,24 +1,76 @@
"""Source-neutral lineage-edge cores.

Primitive-level building blocks for deriving lineage edges. They take only
``task_key`` strings and :class:`~flowx.models.ir.DataAsset` values -- never a
``Pipeline`` or any ``Activity`` -- so the tier-matching, self-edge drop, and
dedup rules live in exactly one place. The discovery-AST derivation in
:mod:`flowx.discovery_lineage` gathers those primitives from a
:class:`~flowx.models.discovery.SourceGraph` and delegates here; the IR-facing
derivation that gathers them from a :class:`~flowx.models.ir.Pipeline` is added
separately in the convert->package lineage work so it stacks on this standard.

They import nothing from ``sources/adf`` or ``sources/airflow``. Everything here
is pure: the functions read their inputs and return new edge lists; nothing is
mutated.
"""Source-neutral lineage derivation over the flowx Pipeline IR.

These functions turn an already-translated :class:`~flowx.models.ir.Pipeline`
into its :class:`~flowx.models.ir.Lineage` block. They operate on IR primitives
only -- ``Activity`` subclasses, ``DataAsset``, ``task_key`` -- and import nothing
from ``sources/adf`` or ``sources/airflow`` so both front-ends share one code
path once they populate ``data_reads`` / ``data_writes`` / ``motif_id``.

The tier-matching, self-edge drop, and dedup rules are factored into two
primitive-level cores -- :func:`control_edges_from_calls` and
:func:`data_edges_from_endpoints` -- that take only ``task_key`` strings and
:class:`DataAsset` values, never a ``Pipeline``. The IR entry points
(:func:`build_control_edges` / :func:`build_data_edges`) gather those primitives
from a pipeline and delegate, and the source-neutral discovery-AST derivation in
:mod:`flowx.discovery_lineage` gathers the same primitives from a
:class:`~flowx.models.discovery.SourceGraph` and delegates too, so both phases
join edges through exactly one implementation.

Everything here is pure: the functions read the pipeline and return new edge
lists / a new :class:`Lineage`; nothing is mutated. :func:`with_lineage` attaches
a block by returning a *new* ``Pipeline`` rather than mutating the input, unlike
the in-place dependency rewrite in ``motifs/collapser.py``.
"""

from __future__ import annotations

from collections.abc import Iterable
import dataclasses
from collections.abc import Iterable, Iterator

from flowx.models.ir import (
Activity,
ControlEdge,
DataAsset,
DataEdge,
ExecutePipelineActivity,
ForEachActivity,
IfConditionActivity,
Lineage,
MotifActivity,
MotifAnnotation,
Pipeline,
RunJobActivity,
SwitchActivity,
)


def walk_activities(activities: list[Activity]) -> Iterator[Activity]:
"""Yield every activity in *activities*, descending into control-flow containers.

Recurses into ForEach inner activities, both If-condition branches, and every
Switch case plus its default branch, so a nested ExecutePipeline or a data
asset buried inside a Switch case is still reached. Motif ``original_activities``
are intentionally not traversed: they are the pre-collapse originals kept for
reference, not live graph members.

from flowx.models.ir import ControlEdge, DataAsset, DataEdge
Args:
activities: Top-level (or already-nested) activity list to walk.

Yields:
Each activity, container nodes included, in depth-first order.
"""
for activity in activities:
yield activity
match activity:
case ForEachActivity():
yield from walk_activities(activity.inner_activities)
case IfConditionActivity():
yield from walk_activities(activity.if_true_activities)
yield from walk_activities(activity.if_false_activities)
case SwitchActivity():
for case_branch in activity.cases:
yield from walk_activities(case_branch.activities)
yield from walk_activities(activity.default_activities)


def control_edges_from_calls(
Expand All @@ -27,10 +79,10 @@ def control_edges_from_calls(
) -> list[ControlEdge]:
"""Assemble deduplicated control edges from raw invocation primitives.

The shared core behind the discovery-AST control-edge derivation (and the
IR-facing derivation added in the convert->package work): it owns the
self-edge drop, the dedup, and the unresolved-callee recording so those rules
live in exactly one place and every phase behaves identically.
The shared core behind :func:`build_control_edges` (IR) and the discovery-AST
control-edge derivation: it owns the self-edge drop, the dedup, and the
unresolved-callee recording so those rules live in exactly one place and both
phases behave identically.

Args:
source_workflow: Name of the calling workflow (pipeline / DAG).
Expand Down Expand Up @@ -65,6 +117,35 @@ def control_edges_from_calls(
return edges


def build_control_edges(pipeline: Pipeline) -> list[ControlEdge]:
"""Derive cross-workflow invocation edges for a pipeline.

Emits one :class:`ControlEdge` per invoking activity -- an
``ExecutePipelineActivity`` (ADF) or a ``RunJobActivity`` (Airflow) -- found
anywhere in the pipeline, including inside ForEach / If / Switch containers
(fan-out is preserved: each call site is its own edge). Edges whose callee
equals the caller are dropped (no self-edges), and identical edges are
collapsed (no duplicates). An unresolved callee is recorded with
``resolved=False`` rather than dropped.

Args:
pipeline: The translated pipeline IR.

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

def _calls() -> Iterator[tuple[str, bool | None, str]]:
for activity in walk_activities(pipeline.tasks):
match activity:
case ExecutePipelineActivity():
yield activity.pipeline_name or "", activity.wait_on_completion, activity.task_key
case RunJobActivity():
yield activity.job_name or "", None, activity.task_key

return control_edges_from_calls(pipeline.name, _calls())


def _match_assets(producer: DataAsset, consumer: DataAsset) -> tuple[str, str, str | None] | None:
"""Decide whether a written asset hands off to a read asset, and how.

Expand Down Expand Up @@ -95,10 +176,9 @@ def data_edges_from_endpoints(
) -> list[DataEdge]:
"""Join producer endpoints to consumer endpoints via the two-tier match.

The shared core behind the discovery-AST data-edge derivation (and the
IR-facing derivation added in the convert->package work): it owns the
:func:`_match_assets` tier logic, the no-self-edge rule, and the dedup, so
every phase joins identically.
The shared core behind :func:`build_data_edges` (IR) and the discovery-AST
data-edge derivation: it owns the :func:`_match_assets` tier logic, the
no-self-edge rule, and the dedup, so both phases join identically.

Args:
producers: ``(task_key, written asset)`` pairs, in first-seen order.
Expand Down Expand Up @@ -137,3 +217,92 @@ def data_edges_from_endpoints(
)
)
return edges


def build_data_edges(pipeline: Pipeline) -> list[DataEdge]:
"""Derive proven producer -> consumer data hand-offs for a pipeline.

A producer is any activity with a ``data_writes`` asset; a consumer any
activity with a ``data_reads`` asset, gathered across the whole pipeline
(ForEach / If / Switch bodies included). Each producer asset is joined against
each consumer asset via :func:`_match_assets`, tagging the edge as an
``identity`` or ``signature`` match. An activity never hands off to itself
(no self-edges), and identical edges are collapsed (no duplicates).

Args:
pipeline: The translated pipeline IR.

Returns:
Deduplicated list of data edges, in first-seen order.
"""
activities = list(walk_activities(pipeline.tasks))
producers = [(activity.task_key, asset) for activity in activities for asset in activity.data_writes]
consumers = [(activity.task_key, asset) for activity in activities for asset in activity.data_reads]
return data_edges_from_endpoints(producers, consumers)


def build_motif_annotations(pipeline: Pipeline) -> list[MotifAnnotation]:
"""Derive motif annotations from the collapsed motif activities in a pipeline.

One annotation per :class:`MotifActivity`, listing the task keys it spans
(the motif task itself plus any member activities that carry the same
``motif_id`` tag). Deduplicated by ``motif_id`` in first-seen order.

Args:
pipeline: The translated pipeline IR.

Returns:
List of motif annotations.
"""
annotations: list[MotifAnnotation] = []
seen: set[str] = set()
activities = list(walk_activities(pipeline.tasks))
for activity in activities:
if not isinstance(activity, MotifActivity):
continue
if activity.motif_id in seen:
continue
seen.add(activity.motif_id)
members = [activity.task_key]
members.extend(
other.task_key for other in activities if other is not activity and other.motif_id == activity.motif_id
)
annotations.append(
MotifAnnotation(
motif_id=activity.motif_id,
member_task_keys=members,
display_name=activity.display_name,
databricks_replacement=activity.databricks_replacement,
notes=list(activity.confidence_notes),
)
)
return annotations


def build_lineage(pipeline: Pipeline) -> Lineage:
"""Compose the full source-neutral lineage block for a pipeline.

Args:
pipeline: The translated pipeline IR.

Returns:
A :class:`Lineage` with control edges, data edges, and motif annotations.
"""
return Lineage(
control_edges=build_control_edges(pipeline),
data_edges=build_data_edges(pipeline),
motifs=build_motif_annotations(pipeline),
)


def with_lineage(pipeline: Pipeline, lineage: Lineage) -> Pipeline:
"""Return a *new* pipeline carrying *lineage*, leaving the input untouched.

Args:
pipeline: The pipeline to copy.
lineage: The lineage block to attach.

Returns:
A shallow copy of *pipeline* with ``lineage`` set.
"""
return dataclasses.replace(pipeline, lineage=lineage)
Loading