From c50e211c36b37b8b0a33ce795ff9fc000769815a Mon Sep 17 00:00:00 2001 From: Matthew Moorcroft Date: Tue, 15 Sep 2026 12:28:25 +0100 Subject: [PATCH] Add convert->package lineage plumbing on top of the discovery standard (#61) PR B of the two-PR split: the CONVERSION/IR side, stacked on the discovery standard (PR A / #66). Adds the IR-Activity and convert->package lineage that PR A deliberately deferred. - models/ir.py: Activity-base fields data_reads / data_writes / motif_id, the MotifActivity base-owned motif_id handling, and Pipeline.lineage. - lineage.py: the IR-Activity-facing layer -- walk_activities, build_control_edges, build_data_edges, build_motif_annotations, build_lineage, with_lineage -- delegating to the neutral cores PR A defines. - ir_serde.py: Pipeline.lineage emission, Activity-base field emission, and the MotifActivity motif_id-via-base serialisation. - bundler/dab_writer.py: full rehydration of the lineage block, Activity-base fields, and MotifActivity through the JSON-reload path. - tests: test_lineage_substrate.py (IR derivation + ir_serde/dab_writer round-trip, incl. the MotifActivity R1 case). Depends on / stacks on PR #66 (PR A). Do not merge before it. Co-authored-by: Isaac --- src/flowx/bundler/dab_writer.py | 66 +++- src/flowx/ir_serde.py | 10 +- src/flowx/lineage.py | 217 +++++++++++-- src/flowx/models/ir.py | 17 +- tests/unit/test_lineage_substrate.py | 454 +++++++++++++++++++++++++++ 5 files changed, 737 insertions(+), 27 deletions(-) create mode 100644 tests/unit/test_lineage_substrate.py diff --git a/src/flowx/bundler/dab_writer.py b/src/flowx/bundler/dab_writer.py index 436900d..8cdf20b 100644 --- a/src/flowx/bundler/dab_writer.py +++ b/src/flowx/bundler/dab_writer.py @@ -30,7 +30,10 @@ from flowx.models.ir import ( Activity, AppendVariableActivity, + ControlEdge, CopyActivity, + DataAsset, + DataEdge, DbtFactoryActivity, DeleteActivity, Dependency, @@ -38,8 +41,10 @@ FilterActivity, ForEachActivity, IfConditionActivity, + Lineage, LookupActivity, MotifActivity, + MotifAnnotation, NotebookActivity, Pipeline, PlaceholderActivity, @@ -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 @@ -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", [])), @@ -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 diff --git a/src/flowx/ir_serde.py b/src/flowx/ir_serde.py index 6783e5b..72b00a6 100644 --- a/src/flowx/ir_serde.py +++ b/src/flowx/ir_serde.py @@ -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 @@ -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) @@ -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 diff --git a/src/flowx/lineage.py b/src/flowx/lineage.py index 9e5d636..447e265 100644 --- a/src/flowx/lineage.py +++ b/src/flowx/lineage.py @@ -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( @@ -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). @@ -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. @@ -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. @@ -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) diff --git a/src/flowx/models/ir.py b/src/flowx/models/ir.py index ffc7e45..b6df8c6 100644 --- a/src/flowx/models/ir.py +++ b/src/flowx/models/ir.py @@ -160,7 +160,8 @@ class MotifAnnotation: Source-neutral record tying a motif id to the tasks that belong to it, so the lineage block can report motifs without depending on how any particular - source detects them. ``member_task_keys`` lists the tasks the motif spans. + source detects them. Each member activity also carries the same + :attr:`Activity.motif_id` tag. Attributes: motif_id: Identifier of the matched motif definition. @@ -234,6 +235,13 @@ class Activity: compute_mode: str | None = None # Collapsed activity_and_notify spec set by the adapter: {destination, events, args, destination_name}. notifications: dict[str, Any] | None = None + # Lineage substrate (#61): source-neutral data assets this activity reads from and writes to, + # populated by per-source extractors in follow-up work. Always lists, never None. + data_reads: list[DataAsset] = field(default_factory=list) + data_writes: list[DataAsset] = field(default_factory=list) + # Id of the motif this activity was folded into (or belongs to); None when it is part of no motif. + # Owned here so every activity type -- not only MotifActivity -- can carry the tag. + motif_id: str | None = None @dataclass(slots=True, kw_only=True) @@ -710,6 +718,10 @@ class PlaceholderActivity(Activity): class MotifActivity(Activity): """Activity produced by collapsing a detected motif pattern. + Redeclares :attr:`Activity.motif_id` as required (the base owns the field so + every activity type can carry the tag and it round-trips through one code + path, but a motif activity always has one). + Attributes: motif_id: Identifier of the matched motif definition. display_name: Human-readable motif name. @@ -764,6 +776,8 @@ class Pipeline: reconciliation_status: Source-audit result for this pipeline. migration_status: Whether the pipeline is included or explicitly excluded. audit: Source-audit counts and transformation ledger. + lineage: Source-neutral lineage block (control/data edges + motif + annotations), or ``None`` when lineage has not been derived. """ name: str @@ -780,6 +794,7 @@ class Pipeline: audit: dict[str, Any] = field(default_factory=dict) translation_configuration: TranslationConfiguration | None = None bundle_variables: dict[str, dict[str, Any]] = field(default_factory=dict) + lineage: Lineage | None = None @dataclass(frozen=True, slots=True) diff --git a/tests/unit/test_lineage_substrate.py b/tests/unit/test_lineage_substrate.py new file mode 100644 index 0000000..1f1d8a8 --- /dev/null +++ b/tests/unit/test_lineage_substrate.py @@ -0,0 +1,454 @@ +"""Unit tests for the source-neutral lineage substrate (#61). + +Covers the pure derivation (:mod:`flowx.lineage`) -- control fan-out, nested and +Switch recursion, the identity-vs-signature join tiers, no self-edges, no +duplicate edges -- and the ``ir_serde`` round-trip for the new IR types plus the +new ``Activity`` base fields, including a MotifActivity round-trip that proves the +R1 ``motif_id`` collision is handled. +""" + +from __future__ import annotations + +import json + +from flowx.bundler.dab_writer import pipeline_dict_to_ir +from flowx.ir_serde import pipeline_to_dict +from flowx.lineage import ( + build_control_edges, + build_data_edges, + build_lineage, + build_motif_annotations, + with_lineage, +) +from flowx.models.ir import ( + ControlEdge, + DataAsset, + DataEdge, + ExecutePipelineActivity, + ForEachActivity, + IfConditionActivity, + Lineage, + MotifActivity, + MotifAnnotation, + NotebookActivity, + Pipeline, + RunJobActivity, + SwitchActivity, + SwitchCase, + WaitActivity, +) + + +def _notebook(task_key: str, *, reads=None, writes=None, motif_id=None) -> NotebookActivity: + return NotebookActivity( + name=task_key, + task_key=task_key, + notebook_path=f"/Shared/{task_key}", + data_reads=list(reads or []), + data_writes=list(writes or []), + motif_id=motif_id, + ) + + +def _execute(task_key: str, callee: str, *, wait: bool = True) -> ExecutePipelineActivity: + return ExecutePipelineActivity(name=task_key, task_key=task_key, pipeline_name=callee, wait_on_completion=wait) + + +# --------------------------------------------------------------------------- # +# Control-edge derivation +# --------------------------------------------------------------------------- # + + +def test_control_edges_fan_out_and_nested_switch_recursion(): + """ExecutePipeline calls are found at top level and inside ForEach/If/Switch.""" + pipeline = Pipeline( + name="parent", + tasks=[ + _execute("call_a", "child_a"), + ForEachActivity( + name="fe", + task_key="fe", + items_expression="@x", + inner_activities=[_execute("call_b", "child_b")], + ), + IfConditionActivity( + name="cond", + task_key="cond", + op="equals", + left="@a", + right="@b", + if_true_activities=[_execute("call_c", "child_c")], + if_false_activities=[_execute("call_d", "child_d")], + ), + SwitchActivity( + name="sw", + task_key="sw", + on_expression="@e", + cases=[SwitchCase(value="one", activities=[_execute("call_e", "child_e")])], + default_activities=[_execute("call_f", "child_f")], + ), + ], + ) + + edges = build_control_edges(pipeline) + + targets = sorted(edge.target_workflow for edge in edges) + assert targets == ["child_a", "child_b", "child_c", "child_d", "child_e", "child_f"] + assert all(edge.source_workflow == "parent" for edge in edges) + # Each call site keeps its own via_task_key (fan-out preserved). + assert {edge.via_task_key for edge in edges} == { + "call_a", + "call_b", + "call_c", + "call_d", + "call_e", + "call_f", + } + + +def test_control_edges_run_job_activity_is_source_neutral(): + """A RunJobActivity (Airflow) produces a control edge just like ExecutePipeline.""" + pipeline = Pipeline( + name="dag_main", + tasks=[RunJobActivity(name="run", task_key="run", job_name="downstream_job")], + ) + + edges = build_control_edges(pipeline) + + assert len(edges) == 1 + assert edges[0].source_workflow == "dag_main" + assert edges[0].target_workflow == "downstream_job" + assert edges[0].via_task_key == "run" + assert edges[0].wait_for_completion is None + assert edges[0].resolved is True + + +def test_control_edges_unresolved_callee_is_recorded_not_dropped(): + """An empty callee is kept with resolved=False rather than silently dropped.""" + pipeline = Pipeline(name="parent", tasks=[_execute("call", "")]) + + edges = build_control_edges(pipeline) + + assert len(edges) == 1 + assert edges[0].target_workflow == "" + assert edges[0].resolved is False + + +def test_control_edges_no_self_edge(): + """A pipeline invoking itself produces no edge.""" + pipeline = Pipeline(name="loop", tasks=[_execute("call", "loop")]) + + assert build_control_edges(pipeline) == [] + + +def test_control_edges_no_duplicate_from_recursion(): + """A single call site nested in a container is emitted exactly once.""" + pipeline = Pipeline( + name="parent", + tasks=[ + ForEachActivity( + name="fe", + task_key="fe", + items_expression="@x", + inner_activities=[_execute("call", "child")], + ) + ], + ) + + edges = build_control_edges(pipeline) + + assert len(edges) == 1 + assert edges[0].target_workflow == "child" + + +# --------------------------------------------------------------------------- # +# Data-edge derivation: identity vs signature tiers +# --------------------------------------------------------------------------- # + + +def test_data_edges_identity_tier_joins_across_different_signatures(): + """Two assets with the same resolved identity match even when their names differ.""" + pipeline = Pipeline( + name="p", + tasks=[ + _notebook("writer", writes=[DataAsset(signature="ds_out", identity="curated.orders")]), + _notebook("reader", reads=[DataAsset(signature="ds_in_other_name", identity="curated.orders")]), + ], + ) + + edges = build_data_edges(pipeline) + + assert len(edges) == 1 + assert edges[0].source_task_key == "writer" + assert edges[0].target_task_key == "reader" + assert edges[0].match_kind == "identity" + assert edges[0].match_key == "curated.orders" + assert edges[0].identity == "curated.orders" + + +def test_data_edges_signature_tier_when_identity_unresolved(): + """When identity is unresolvable, matching falls back to the neutral signature.""" + pipeline = Pipeline( + name="p", + tasks=[ + _notebook("writer", writes=[DataAsset(signature="shared_ds")]), + _notebook("reader", reads=[DataAsset(signature="shared_ds")]), + ], + ) + + edges = build_data_edges(pipeline) + + assert len(edges) == 1 + assert edges[0].match_kind == "signature" + assert edges[0].match_key == "shared_ds" + assert edges[0].identity is None + + +def test_data_edges_distinct_identities_do_not_fall_back_to_signature(): + """Two resolved-but-different identities never manufacture a signature edge (#36).""" + pipeline = Pipeline( + name="p", + tasks=[ + _notebook("writer", writes=[DataAsset(signature="shared", identity="a.first")]), + _notebook("reader", reads=[DataAsset(signature="shared", identity="b.second")]), + ], + ) + + assert build_data_edges(pipeline) == [] + + +def test_data_edges_fan_out_one_writer_many_readers(): + """One producer handing off to several consumers yields one edge each.""" + pipeline = Pipeline( + name="p", + tasks=[ + _notebook("writer", writes=[DataAsset(signature="ds", identity="x.y")]), + _notebook("reader_one", reads=[DataAsset(signature="ds", identity="x.y")]), + _notebook("reader_two", reads=[DataAsset(signature="ds", identity="x.y")]), + ], + ) + + edges = build_data_edges(pipeline) + + assert sorted(edge.target_task_key for edge in edges) == ["reader_one", "reader_two"] + assert all(edge.source_task_key == "writer" for edge in edges) + + +def test_data_edges_no_self_edge(): + """An activity that both writes and reads the same asset does not edge to itself.""" + pipeline = Pipeline( + name="p", + tasks=[ + _notebook( + "roundtrip", + writes=[DataAsset(signature="ds", identity="x.y")], + reads=[DataAsset(signature="ds", identity="x.y")], + ) + ], + ) + + assert build_data_edges(pipeline) == [] + + +def test_data_edges_no_duplicate_from_repeated_asset(): + """A producer listing the same asset twice still yields a single edge.""" + pipeline = Pipeline( + name="p", + tasks=[ + _notebook( + "writer", + writes=[DataAsset(signature="ds", identity="x.y"), DataAsset(signature="ds", identity="x.y")], + ), + _notebook("reader", reads=[DataAsset(signature="ds", identity="x.y")]), + ], + ) + + edges = build_data_edges(pipeline) + + assert len(edges) == 1 + + +def test_data_edges_nested_switch_recursion(): + """A producer buried in a Switch case hands off to a top-level consumer.""" + pipeline = Pipeline( + name="p", + tasks=[ + SwitchActivity( + name="sw", + task_key="sw", + on_expression="@e", + cases=[ + SwitchCase( + value="one", + activities=[_notebook("writer", writes=[DataAsset(signature="ds", identity="x.y")])], + ) + ], + default_activities=[], + ), + _notebook("reader", reads=[DataAsset(signature="ds", identity="x.y")]), + ], + ) + + edges = build_data_edges(pipeline) + + assert len(edges) == 1 + assert edges[0].source_task_key == "writer" + assert edges[0].target_task_key == "reader" + + +# --------------------------------------------------------------------------- # +# Motif annotations + composition + purity +# --------------------------------------------------------------------------- # + + +def test_build_motif_annotations_groups_members_by_tag(): + """A MotifActivity plus tagged members become one annotation over their task keys.""" + pipeline = Pipeline( + name="p", + tasks=[ + MotifActivity( + name="motif", + task_key="motif_auto_loader", + motif_id="auto_loader", + display_name="Auto Loader", + databricks_replacement="auto_loader", + matched_activity_names=["Copy A", "Copy B"], + confidence_notes=["matched on file source"], + ), + _notebook("member", motif_id="auto_loader"), + _notebook("unrelated"), + ], + ) + + annotations = build_motif_annotations(pipeline) + + assert len(annotations) == 1 + assert annotations[0].motif_id == "auto_loader" + assert annotations[0].member_task_keys == ["motif_auto_loader", "member"] + assert annotations[0].display_name == "Auto Loader" + assert annotations[0].databricks_replacement == "auto_loader" + + +def test_with_lineage_is_pure(): + """with_lineage returns a new pipeline and never mutates the input.""" + pipeline = Pipeline(name="p", tasks=[_notebook("n")]) + lineage = build_lineage(pipeline) + + updated = with_lineage(pipeline, lineage) + + assert pipeline.lineage is None + assert updated is not pipeline + assert updated.lineage is lineage + + +# --------------------------------------------------------------------------- # +# ir_serde round-trips +# --------------------------------------------------------------------------- # + + +def test_serde_round_trip_new_activity_fields_and_lineage_block(): + """data_reads/data_writes/motif_id and the lineage block survive JSON round-trip.""" + pipeline = Pipeline( + name="p", + tasks=[ + _notebook( + "writer", + writes=[DataAsset(signature="ds_out", identity="curated.orders", asset_type="table")], + motif_id="auto_loader", + ), + _notebook( + "reader", + reads=[ + DataAsset( + signature="ds_in", + identity="curated.orders", + asset_type="table", + properties={"format": "delta"}, + ) + ], + ), + WaitActivity(name="pause", task_key="pause", wait_time_seconds=5), + ], + lineage=Lineage( + control_edges=[ + ControlEdge( + source_workflow="p", + target_workflow="child", + via_task_key="writer", + wait_for_completion=True, + resolved=True, + ) + ], + data_edges=[ + DataEdge( + source_task_key="writer", + target_task_key="reader", + match_kind="identity", + match_key="curated.orders", + identity="curated.orders", + asset_type="table", + ) + ], + motifs=[ + MotifAnnotation( + motif_id="auto_loader", + member_task_keys=["writer"], + display_name="Auto Loader", + databricks_replacement="auto_loader", + notes=["note"], + ) + ], + ), + ) + + reloaded, _ = pipeline_dict_to_ir(json.loads(json.dumps(pipeline_to_dict(pipeline)))) + + writer = reloaded.tasks[0] + reader = reloaded.tasks[1] + assert writer.motif_id == "auto_loader" + assert writer.data_writes == [DataAsset(signature="ds_out", identity="curated.orders", asset_type="table")] + assert reader.data_reads == [ + DataAsset(signature="ds_in", identity="curated.orders", asset_type="table", properties={"format": "delta"}) + ] + # A task without lineage fields rehydrates to empty lists / None, never missing. + assert reloaded.tasks[2].data_reads == [] + assert reloaded.tasks[2].data_writes == [] + assert reloaded.tasks[2].motif_id is None + + assert reloaded.lineage == pipeline.lineage + + +def test_serde_round_trip_motif_activity_no_kwarg_collision_r1(): + """A MotifActivity round-trips without the R1 'multiple values for motif_id' TypeError.""" + pipeline = Pipeline( + name="p", + tasks=[ + MotifActivity( + name="motif", + task_key="motif_auto_loader", + motif_id="auto_loader", + display_name="Auto Loader", + databricks_replacement="auto_loader", + matched_activity_names=["Copy A", "Copy B"], + data_reads=[DataAsset(signature="src", identity="raw.src")], + ) + ], + ) + + reloaded, _ = pipeline_dict_to_ir(json.loads(json.dumps(pipeline_to_dict(pipeline)))) + + task = reloaded.tasks[0] + assert isinstance(task, MotifActivity) + assert task.motif_id == "auto_loader" + assert task.display_name == "Auto Loader" + assert task.matched_activity_names == ["Copy A", "Copy B"] + assert task.data_reads == [DataAsset(signature="src", identity="raw.src")] + + +def test_lineage_block_always_emits_lists_never_null(): + """An attached empty lineage serialises its edge collections as lists, not null.""" + pipeline = with_lineage(Pipeline(name="p", tasks=[_notebook("n")]), Lineage()) + + serialised = pipeline_to_dict(pipeline) + + assert serialised["lineage"] == {"control_edges": [], "data_edges": [], "motifs": []}