diff --git a/src/flowx/discovery_lineage.py b/src/flowx/discovery_lineage.py new file mode 100644 index 0000000..1e9d0c0 --- /dev/null +++ b/src/flowx/discovery_lineage.py @@ -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)) diff --git a/src/flowx/discovery_serde.py b/src/flowx/discovery_serde.py index 9dc77fb..ee7a881 100644 --- a/src/flowx/discovery_serde.py +++ b/src/flowx/discovery_serde.py @@ -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 @@ -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, @@ -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]: @@ -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: @@ -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", ""), @@ -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: diff --git a/src/flowx/lineage.py b/src/flowx/lineage.py index a417718..447e265 100644 --- a/src/flowx/lineage.py +++ b/src/flowx/lineage.py @@ -6,6 +6,16 @@ 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 @@ -15,7 +25,7 @@ from __future__ import annotations import dataclasses -from collections.abc import Iterator +from collections.abc import Iterable, Iterator from flowx.models.ir import ( Activity, @@ -63,49 +73,43 @@ def walk_activities(activities: list[Activity]) -> Iterator[Activity]: yield from walk_activities(activity.default_activities) -def build_control_edges(pipeline: Pipeline) -> list[ControlEdge]: - """Derive cross-workflow invocation edges for a pipeline. +def control_edges_from_calls( + source_workflow: str, + calls: Iterable[tuple[str, bool | None, str]], +) -> list[ControlEdge]: + """Assemble deduplicated control edges from raw invocation primitives. - 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. + 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: - pipeline: The translated pipeline IR. + source_workflow: Name of the calling workflow (pipeline / DAG). + calls: One ``(target_workflow, wait_for_completion, via_task_key)`` triple + per call site, in the order they should be considered. ``target_workflow`` + may be empty when the callee could not be resolved from a partial export. Returns: - Deduplicated list of control edges, in first-seen order. + Deduplicated control edges in first-seen order. A call whose target equals + ``source_workflow`` is dropped (no self-edge); an empty target is kept with + ``resolved=False`` rather than dropped. """ edges: list[ControlEdge] = [] seen: set[tuple[str, str, str]] = set() - for activity in walk_activities(pipeline.tasks): - target: str | None - wait: bool | None - match activity: - case ExecutePipelineActivity(): - target = activity.pipeline_name - wait = activity.wait_on_completion - case RunJobActivity(): - target = activity.job_name - wait = None - case _: - continue + for target, wait, via_task_key in calls: target_name = target or "" - if target_name and target_name == pipeline.name: + if target_name and target_name == source_workflow: continue - key = (pipeline.name, target_name, activity.task_key) + key = (source_workflow, target_name, via_task_key) if key in seen: continue seen.add(key) edges.append( ControlEdge( - source_workflow=pipeline.name, + source_workflow=source_workflow, target_workflow=target_name, - via_task_key=activity.task_key, + via_task_key=via_task_key, wait_for_completion=wait, resolved=bool(target_name), ) @@ -113,6 +117,35 @@ def build_control_edges(pipeline: Pipeline) -> list[ControlEdge]: 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. @@ -137,30 +170,32 @@ def _match_assets(producer: DataAsset, consumer: DataAsset) -> tuple[str, str, s return None -def build_data_edges(pipeline: Pipeline) -> list[DataEdge]: - """Derive proven producer -> consumer data hand-offs for a pipeline. +def data_edges_from_endpoints( + producers: Iterable[tuple[str, DataAsset]], + consumers: Iterable[tuple[str, DataAsset]], +) -> list[DataEdge]: + """Join producer endpoints to consumer endpoints via the two-tier match. - 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). + 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: - pipeline: The translated pipeline IR. + producers: ``(task_key, written asset)`` pairs, in first-seen order. + consumers: ``(task_key, read asset)`` pairs, in first-seen order. Returns: - Deduplicated list of data edges, in first-seen order. + Deduplicated data edges in first-seen order. A producer never hands off to + a consumer sharing its ``task_key`` (no self-edge), and identical + ``(producer, consumer, match_kind, match_key)`` edges are collapsed. """ - 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] + producer_list = list(producers) + consumer_list = list(consumers) edges: list[DataEdge] = [] seen: set[tuple[str, str, str, str]] = set() - for producer_key, producer_asset in producers: - for consumer_key, consumer_asset in consumers: + for producer_key, producer_asset in producer_list: + for consumer_key, consumer_asset in consumer_list: if producer_key == consumer_key: continue matched = _match_assets(producer_asset, consumer_asset) @@ -184,6 +219,28 @@ def build_data_edges(pipeline: Pipeline) -> list[DataEdge]: 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. diff --git a/src/flowx/models/discovery.py b/src/flowx/models/discovery.py index 1590c0b..971d89f 100644 --- a/src/flowx/models/discovery.py +++ b/src/flowx/models/discovery.py @@ -35,7 +35,7 @@ from dataclasses import dataclass, field from typing import Any -from flowx.models.ir import DataAsset +from flowx.models.ir import DataAsset, Lineage # --------------------------------------------------------------------------- # # Well-known discriminators and concepts (open vocabularies). @@ -248,6 +248,11 @@ class SourceGraph: tags: Free-form label list (ADF annotations, Airflow user tags). tasks: Top-level nodes; control flow nests further nodes via :class:`ContainerNode`. + lineage: Source-neutral lineage block (control/data edges) derived over + this graph, or ``None`` when lineage has not been derived. Reuses the + #61 :class:`~flowx.models.ir.Lineage` type so the discovery AST and + the IR lineage substrate share one vocabulary; populated by + :mod:`flowx.discovery_lineage` in the discover phase. properties: Free-form bag for graph-level platform-specific attributes. extensions: Alias-free overflow for anything else source-specific (e.g. ADF ``folder``, Airflow global Variables). @@ -262,6 +267,7 @@ class SourceGraph: schedule: ScheduleSpec | None = None tags: list[str] = field(default_factory=list) tasks: list[SourceNode] = field(default_factory=list) + lineage: Lineage | None = None properties: dict[str, Any] = field(default_factory=dict) extensions: dict[str, Any] = field(default_factory=dict) raw: dict[str, Any] | None = None diff --git a/src/flowx/sources/adf/dataset_lineage.py b/src/flowx/sources/adf/dataset_lineage.py new file mode 100644 index 0000000..ba51aa7 --- /dev/null +++ b/src/flowx/sources/adf/dataset_lineage.py @@ -0,0 +1,495 @@ +"""Resolve ADF dataset references into source-neutral :class:`DataAsset` values. + +This is the ADF half of data-lineage population (#62b). It re-homes the +dataset-identity resolver and the path-signature logic first written for the +closed #36 ADF-only lineage attempt, but emits the shared #61 two-tier +:class:`~flowx.models.ir.DataAsset` (``identity`` + ``signature``) instead of a +bespoke edge type, so the source-neutral join in :mod:`flowx.lineage` / +:mod:`flowx.discovery_lineage` does the matching. + +Two tiers, exactly as #36 established them: + +* **identity** -- the resolved physical location of the asset (``schema.table`` or + a concrete ``abfss://`` path). Present only when it resolves *deterministically* + from literals; a parameterised reference is never guessed at and leaves + ``identity`` unset. This is the strong join key. +* **signature** -- the *path-derived* weak key. For an asset with a resolved + identity the signature mirrors that identity, so the weak tier can never join a + resolved asset to an unresolved one on a coincidence. For an unresolved reference + the signature is the *structural path signature* (#36's "expression" tier: the + literal path skeleton plus its parameter-slot count) when the path has a literal + anchor. It is **never** the bare dataset reference name: two unrelated opaque + references that merely share a name must not join (#36's explicit rule), so when + neither a physical identity nor a path-anchored signature is available the + signature is left empty and the asset cannot participate in signature matching. + +Only literal, provable values ever become an ``identity`` -- the resolver returns +``None`` rather than guessing, which is what stopped #36's spurious edges. +""" + +from __future__ import annotations + +import json +import re +from collections.abc import Iterator +from typing import Any + +from flowx.models.adf_ast import AdfActivity, AdfDatasetReference, AdfDefinitions +from flowx.models.ir import DataAsset, TranslationContext +from flowx.parser.expression_parser import resolve_expression, resolve_interpolated_string + +_ACCOUNT_NAME_RE = re.compile(r"AccountName=([A-Za-z0-9]+)", re.IGNORECASE) +_DATASET_PARAM_RE = re.compile(r"^@dataset\(\)\.([A-Za-z_][A-Za-z0-9_]*)$") + +# Runtime references inside a path expression. Each is a value only knowable at +# run time; for a *structural* signature we collapse them all to one slot token so +# that, e.g., a writer's ``pipeline().parameters.entityID`` and a reader's +# ``item().entityID`` (the same value passed down a ForEach) share the same shape. +_PARAM_REF_RE = re.compile( + r"pipeline\(\)\.parameters\.\w+" + r"|item\(\)(?:\.\w+)*" + r"|variables\('[^']*'\)" + r"|dataset\(\)\.\w+" + r"|activity\('[^']*'\)\.[\w.]+" +) +_QUOTED_LITERAL_RE = re.compile(r"'([^']*)'") + + +# --------------------------------------------------------------------------- # +# Public entry point +# --------------------------------------------------------------------------- # + + +def activity_data_assets( + activity: AdfActivity, + definitions: AdfDefinitions, + context: TranslationContext | None = None, +) -> tuple[list[DataAsset], list[DataAsset]]: + """Resolve the assets an ADF activity reads and writes. + + Captures **every** input and output, not just index 0: a Copy reads its + ``source`` dataset(s) and writes its ``sink`` dataset(s), a Lookup reads its + ``dataset``, and any activity-level ``inputs`` / ``outputs`` slots are all + included. A dataset named in both an activity-level slot and ``typeProperties`` + is counted once per side so an activity does not emit two identical assets. + + Args: + activity: The ADF activity to resolve. + definitions: All loaded ADF definitions (datasets + linked services), + needed to resolve a reference to its physical identity. + context: Optional translation context for expression resolution; a default + (empty) context is used when none is given, mirroring the deterministic + discover-time resolution the closed #36 attempt used. + + Returns: + ``(data_reads, data_writes)`` as lists of :class:`DataAsset`. + """ + resolution_context = context if context is not None else TranslationContext() + reads = [ + _dataset_ref_to_asset(reference, definitions, resolution_context) + for reference in _activity_dataset_refs(activity, produced=False) + ] + writes = [ + _dataset_ref_to_asset(reference, definitions, resolution_context) + for reference in _activity_dataset_refs(activity, produced=True) + ] + return reads, writes + + +def _dataset_ref_to_asset( + dataset_ref: AdfDatasetReference, + definitions: AdfDefinitions, + context: TranslationContext, +) -> DataAsset: + """Turn one dataset reference into a two-tier :class:`DataAsset`. + + Signature is derived only from the resolved *physical* location -- the identity + when it resolves, else the structural path signature. It is **never** the bare + dataset reference name (#36's hard rule): two unrelated opaque references that + merely share a name must not join, so when neither a physical identity nor a + path-anchored signature is available the signature is left empty. An empty + signature is falsy, so :func:`~flowx.lineage._match_assets` cannot use it as a + join key -- the asset is still captured as a read / write for reporting, it just + cannot manufacture a signature-tier edge. + """ + identity = resolve_dataset_identity(dataset_ref, definitions, context) + if identity is not None: + # Mirror the identity into the signature so the weak tier never joins a + # resolved asset to an unresolved one that merely shares a physical value. + return DataAsset(signature=identity, identity=identity, asset_type=_asset_type(dataset_ref, definitions)) + path_signature = _path_signature(dataset_ref.parameters) + return DataAsset( + signature=path_signature if path_signature is not None else "", + identity=None, + asset_type=_asset_type(dataset_ref, definitions), + ) + + +# --------------------------------------------------------------------------- # +# Dataset reference gathering (all inputs / outputs) +# --------------------------------------------------------------------------- # + + +def _typeprops_dataset_ref(candidate: object) -> AdfDatasetReference | None: + """Build a dataset reference from a ``typeProperties`` source/sink/dataset slot. + + Carries the slot's ``parameters`` (Lookup / Delete / GetMetadata put the + dataset call-site params here) so the path signature can be computed. + """ + if isinstance(candidate, dict): + name = candidate.get("referenceName") + if isinstance(name, str) and name: + parameters = candidate.get("parameters") + return AdfDatasetReference( + reference_name=name, + parameters=parameters if isinstance(parameters, dict) else None, + ) + return None + + +def _activity_dataset_refs(activity: AdfActivity, *, produced: bool) -> Iterator[AdfDatasetReference]: + """Yield the dataset references an activity writes (produced) or reads (not). + + Gathers activity-level ``inputs`` / ``outputs`` **and** the ``typeProperties`` + ``source`` / ``sink`` / ``dataset`` slots. De-duplication is by + ``(reference_name, parameter binding)``, not by name alone: a dataset named in + both an activity slot and ``typeProperties`` with the *same* call-site params is + the same physical asset and is yielded once, but two uses of the *same* + parameterised dataset with *different* params (``ds(tbl=orders)`` vs + ``ds(tbl=customers)``) resolve to distinct physical assets and are both kept. + """ + type_properties = activity.type_properties or {} + candidates: list[AdfDatasetReference] = [] + if produced: + candidates.extend(activity.outputs or []) + sink_reference = _typeprops_dataset_ref(type_properties.get("sink")) + if sink_reference is not None: + candidates.append(sink_reference) + else: + candidates.extend(activity.inputs or []) + for key in ("source", "dataset"): + read_reference = _typeprops_dataset_ref(type_properties.get(key)) + if read_reference is not None: + candidates.append(read_reference) + + seen: set[tuple[str, str]] = set() + for reference in candidates: + dedupe_key = (reference.reference_name, _parameter_binding_key(reference)) + if dedupe_key in seen: + continue + seen.add(dedupe_key) + yield reference + + +def _parameter_binding_key(reference: AdfDatasetReference) -> str: + """Stable key for a reference's call-site parameter binding. + + Two references with the same name collapse only when their parameters match, so + distinct bindings that resolve to distinct physical assets survive. Sorted keys + make the string order-independent; ``default=str`` keeps it total for any value + an ADF export can carry. + """ + return json.dumps(reference.parameters or {}, sort_keys=True, default=str) + + +# --------------------------------------------------------------------------- # +# Physical identity resolution (tier 1) +# --------------------------------------------------------------------------- # + + +def resolve_dataset_identity( + dataset_ref: AdfDatasetReference, + definitions: AdfDefinitions, + context: TranslationContext | None = None, +) -> str | None: + """Deterministic physical identity for a dataset reference. + + Returns ``"schema.table"`` when a table is resolvable, else a storage path, + else ``None`` (never a guess). Used to join producers to consumers on the same + physical asset even when their ADF dataset names differ. + + Parameterised values (ADF expressions or DAB-ref placeholders) are treated as + unresolvable and return ``None`` -- they must never be used as identity keys + because two unrelated pipelines sharing a parameter name would collide on the + same placeholder string. + """ + resolution_context = context if context is not None else TranslationContext() + properties = _dataset_props(dataset_ref, definitions) + if properties is None: + return None + schema, table = _resolve_table_reference(dataset_ref, properties, resolution_context) + if table: + identity = f"{schema}.{table}" if schema else table + return identity if _is_physical(identity) else None + path = _resolve_dataset_path(properties, definitions) + return path if (path and _is_physical(path)) else None + + +def _dataset_props(dataset_ref: AdfDatasetReference, definitions: AdfDefinitions) -> dict[str, Any] | None: + """Return the ``properties`` dict for a dataset reference, or ``None``.""" + dataset = definitions.get_dataset(dataset_ref.reference_name) + if not dataset: + return None + return dict(dataset.properties or {}) + + +def _resolve_param_value( + raw: Any, + dataset_params: dict[str, Any], + context: TranslationContext, +) -> str: + """Resolve a single ADF location / table field to a string.""" + if raw is None: + return "" + if isinstance(raw, dict) and raw.get("type") == "Expression": + raw = raw.get("value", "") + if isinstance(raw, (list, dict)): + return "" + if not isinstance(raw, str): + return str(raw) + text = raw + + match = _DATASET_PARAM_RE.match(text.strip()) + if match: + parameter_name = match.group(1) + return _resolve_param_value(dataset_params.get(parameter_name, ""), dataset_params, context) + + if "@{" in text: + return resolve_interpolated_string(text, context) + + if text.startswith("@"): + result = resolve_expression(text, context) + if result is not None and result.kind in ("literal", "dab_ref"): + return result.value + return text + + return text + + +def _effective_dataset_params(dataset_ref: AdfDatasetReference, dataset_props: dict[str, Any]) -> dict[str, Any]: + """Effective parameter map: declared dataset defaults first, call-site overrides win.""" + declared = dataset_props.get("parameters") or {} + effective: dict[str, Any] = {} + for name, spec in declared.items(): + if isinstance(spec, dict) and "defaultValue" in spec: + effective[name] = spec["defaultValue"] + if dataset_ref.parameters: + effective.update(dict(dataset_ref.parameters)) + return effective + + +def _resolve_table_reference( + dataset_ref: AdfDatasetReference, + dataset_props: dict[str, Any] | None, + context: TranslationContext, +) -> tuple[str | None, str | None]: + """Resolve ``(schema, table)`` from a dataset reference. + + Handles both the nested ``typeProperties`` shape and the + ``schemaTypePropertiesSchema`` flattened form ``az datafactory dataset show`` + emits. ADF parameter expressions resolve against the reference's effective + parameter map. + """ + if not dataset_props: + return None, None + type_props = dataset_props.get("typeProperties") if isinstance(dataset_props.get("typeProperties"), dict) else None + effective_params = _effective_dataset_params(dataset_ref, dataset_props) + schema_raw = _pick_dataset_field( + type_props, + dataset_props, + ("schema", "database"), + ("schemaTypePropertiesSchema", "database"), + ) + table_raw = _pick_dataset_field( + type_props, + dataset_props, + ("table", "tableName"), + ("table", "tableName"), + ) + schema = _resolve_param_value(schema_raw, effective_params, context) if schema_raw is not None else None + table = _resolve_param_value(table_raw, effective_params, context) if table_raw is not None else None + return (schema or None), (table or None) + + +def _pick_dataset_field( + type_props: dict[str, Any] | None, + dataset_props: dict[str, Any], + nested_keys: tuple[str, ...], + flat_keys: tuple[str, ...], +) -> Any: + """First populated dataset field across the nested and az-flattened shapes. + + Empty strings, empty lists, and ``None`` are skipped so a column-schema + artifact like ``schema: []`` does not shadow the real database schema stored + under a flattened key. + """ + candidates: list[Any] = [] + if type_props is not None: + candidates.extend(type_props.get(key) for key in nested_keys) + candidates.extend(dataset_props.get(key) for key in flat_keys) + for value in candidates: + if value is None: + continue + if isinstance(value, (list, dict)) and not value: + continue + if isinstance(value, str) and not value.strip(): + continue + return value + return None + + +def _resolve_dataset_path(dataset_props: dict[str, Any], definitions: AdfDefinitions) -> str | None: + """Resolve a dataset's storage path from its location + backing linked service.""" + type_props = dataset_props.get("typeProperties") or dataset_props + location = type_props.get("location") or {} + if not isinstance(location, dict): + return None + + file_system = location.get("fileSystem") or location.get("container") or "" + folder_path = location.get("folderPath") or "" + if isinstance(file_system, dict) or isinstance(folder_path, dict): + return None # parameterised location; not a deterministic identity + + linked_service_reference = dataset_props.get("linkedServiceName") or {} + if isinstance(linked_service_reference, dict): + linked_service_name = linked_service_reference.get("referenceName", "") + else: + linked_service_name = str(linked_service_reference) + linked_service = definitions.get_linked_service(linked_service_name) if linked_service_name else None + account = _resolve_storage_account(linked_service) + if not account: + return None + + return f"abfss://{file_system}@{account}.dfs.core.windows.net/{folder_path}".rstrip("/") + + +def _resolve_storage_account(linked_service: Any) -> str | None: + """Pull a storage account name out of a linked service, if present.""" + if linked_service is None: + return None + type_props = linked_service.properties.get("typeProperties") or linked_service.properties + + url = type_props.get("url") or "" + if isinstance(url, str) and url: + host = url.replace("https://", "").split("/", 1)[0] + host_no_port = host.split(":", 1)[0] + if "." in host_no_port: + return host_no_port.split(".", 1)[0] + + sas_uri = type_props.get("sasUri") or "" + if isinstance(sas_uri, str) and sas_uri: + host = sas_uri.split("?", 1)[0].replace("https://", "").split("/", 1)[0] + if "." in host: + return host.split(".", 1)[0] + + # Plaintext connection string (rare in az exports -- usually masked). + connection_string = type_props.get("connectionString") + if isinstance(connection_string, str): + match = _ACCOUNT_NAME_RE.search(connection_string) + if match: + return match.group(1) + if isinstance(connection_string, dict): + value = connection_string.get("value", "") + match = _ACCOUNT_NAME_RE.search(value) + if match: + return match.group(1) + + # AWS -- bucket name lives on the dataset, account is implicit; nothing useful + # to return at the linked-service level for S3 / GCS. + return None + + +def _is_physical(value: str) -> bool: + """Return ``True`` only when *value* is a literal (physical) identifier. + + A value is NOT physical when it still contains an unresolved marker -- a + DAB-ref placeholder (``{{`` ... ``}}``), a leftover ADF interpolation + fragment (``@{``), or a bare ADF expression (starts with ``@``). + """ + stripped = value.lstrip() + return not ("{{" in value or "@{" in value or stripped.startswith("@")) + + +# --------------------------------------------------------------------------- # +# Structural path signature (tier 2, #36's "expression" tier) +# --------------------------------------------------------------------------- # + + +def _path_signature(parameters: dict[str, Any] | None) -> str | None: + """Structural signature of a reference's parameterised folderPath / fileName. + + Requires a resolvable **folderPath** literal anchor. A file name alone is too + weak a discriminator: many unrelated activities write ``.csv`` / ``.json`` files + to opaque parameterised folders, so a signature built only from a file extension + would join them all -- re-creating the explosion the identity-only join avoids. + Anchoring on the literal folder segment keeps the match specific to a real, + named location. + + ``None`` when the folder path has no literal segment to anchor on. + """ + if not parameters: + return None + folder_signature = _normalize_path_expression(parameters.get("folderPath")) + if folder_signature is None: + return None + file_signature = _normalize_path_expression(parameters.get("fileName")) + return f"FP[{folder_signature}]/FN[{file_signature}]" + + +def _normalize_path_expression(expression: Any) -> str | None: + """Reduce a (possibly parameterised) ADF path expression to a structural signature. + + Keeps the literal path segments and collapses every runtime reference to a + single ``
`` slot, so the result captures the path *shape* (literal skeleton + plus slot count) without guessing the runtime value. Returns ``None`` when + there is no literal segment to anchor on (a signature of only slots is too weak + a join key -- never guess). + """ + if isinstance(expression, dict): + expression = expression.get("value", "") + if not isinstance(expression, str) or not expression.strip(): + return None + text = expression.strip() + if "@" not in text: + # A bare literal value (no ADF expression): the whole string is the literal + # path / filename, with no runtime slots. + literal = re.sub(r"/+", "/", text).strip("/") + return f"{literal}|slots=0" if literal else None + marked = _PARAM_REF_RE.sub("
", text) + literal = re.sub(r"/+", "/", "".join(_QUOTED_LITERAL_RE.findall(marked))).strip("/") + if not literal: + return None + return f"{literal}|slots={marked.count('
')}" + + +# --------------------------------------------------------------------------- # +# Asset-type classification (best-effort neutral kind) +# --------------------------------------------------------------------------- # + +_TABLE_HINTS = ("table", "sql", "database") +_FILE_HINTS = ("delimited", "parquet", "orc", "avro", "json", "binary", "excel", "xml", "blob", "adls", "file") + + +def _asset_type(dataset_ref: AdfDatasetReference, definitions: AdfDefinitions) -> str | None: + """Best-effort neutral asset kind (``"table"`` / ``"file"``), or ``None``. + + Prefers the shape of the dataset's typeProperties (a ``location`` block means a + file, a ``table`` / ``schema`` means a table) and falls back to keyword hints in + the dataset's ADF type string. Returns ``None`` when nothing is conclusive + rather than guessing. + """ + dataset = definitions.get_dataset(dataset_ref.reference_name) + if dataset is None: + return None + type_props = dataset.properties.get("typeProperties") + if isinstance(type_props, dict): + if isinstance(type_props.get("location"), dict): + return "file" + if type_props.get("table") or type_props.get("tableName") or type_props.get("schema"): + return "table" + dataset_type = (dataset.type or "").lower() + if any(hint in dataset_type for hint in _TABLE_HINTS): + return "table" + if any(hint in dataset_type for hint in _FILE_HINTS): + return "file" + return None diff --git a/src/flowx/sources/adf/discovery_mapping.py b/src/flowx/sources/adf/discovery_mapping.py index 5841492..0555a31 100644 --- a/src/flowx/sources/adf/discovery_mapping.py +++ b/src/flowx/sources/adf/discovery_mapping.py @@ -15,10 +15,17 @@ the target-side translation strategy -- rides in the ``properties`` / ``extensions`` seams rather than being dropped or forced into a field. -Populating a node's ``data_reads`` / ``data_writes`` from ADF dataset resolvers, -deriving control edges, and Switch-aware lineage walking are intentionally **not** -done here -- they belong to the follow-up slice (#62b). Nodes therefore carry -empty read/write lists for now. +Data lineage (#62b) is populated here now: every node's ``data_reads`` / +``data_writes`` are resolved from the ADF dataset references via +:func:`~flowx.sources.adf.dataset_lineage.activity_data_assets` (the two-tier +identity / signature model), and every ``ExecutePipeline`` node records the child +pipeline it invokes under the neutral +:data:`~flowx.discovery_lineage.INVOKES_WORKFLOW_PROPERTY` marker. Because +:func:`_activity_to_node` recurses into every control-flow branch, this population +is Switch-aware for free -- a Copy or ExecutePipeline nested inside a Switch case +(or default), a ForEach / Until body, or either If branch carries its assets and +marker like any top-level node. The graph's :attr:`SourceGraph.lineage` block is +then derived by the source-neutral :func:`~flowx.discovery_lineage.build_graph_lineage`. The control-flow container shape follows :class:`~flowx.models.discovery.ContainerNode`: an ``IfCondition`` becomes ``{"true": [...], "false": [...]}``, a ``ForEach`` / @@ -32,6 +39,7 @@ from typing import Any from flowx.discovery_inventory import STRATEGY_PROPERTY +from flowx.discovery_lineage import INVOKES_WAIT_PROPERTY, INVOKES_WORKFLOW_PROPERTY, build_graph_lineage from flowx.models.adf_ast import AdfActivity, AdfDefinitions, AdfParameter, AdfPipeline, AdfTrigger, AdfVariable from flowx.models.discovery import ( CONCEPT_BRANCH, @@ -55,6 +63,7 @@ SourceGraph, SourceNode, ) +from flowx.sources.adf.dataset_lineage import activity_data_assets from flowx.sources.adf.loader import classify_activity # Neutral trigger category for each ADF trigger type. Anything unrecognised maps @@ -96,10 +105,19 @@ def adf_definitions_to_source_graphs(definitions: AdfDefinitions) -> list[Source Pipeline order is preserved so downstream consumers see the same ordering the loader produced. ADF triggers are mapped into :class:`ScheduleSpec` and attached to each pipeline they reference (see :func:`_attach_schedules`), so schedule / - trigger information is preserved rather than dropped. + trigger information is preserved rather than dropped. Each graph's + :attr:`SourceGraph.lineage` block is derived (control + data edges) once its + nodes' reads / writes and invocation markers are populated -- data edges are + joined within a single graph, matching the per-workflow scope of the shared + :class:`~flowx.models.ir.Lineage` model. + + ``definitions`` is threaded down to the activity mapper so a node's data assets + resolve against the factory's datasets and linked services. """ - graphs = [adf_pipeline_to_source_graph(pipeline) for pipeline in definitions.pipelines] + graphs = [adf_pipeline_to_source_graph(pipeline, definitions) for pipeline in definitions.pipelines] _attach_schedules(graphs, definitions.triggers) + for graph in graphs: + graph.lineage = build_graph_lineage(graph) return graphs @@ -182,8 +200,17 @@ def _trigger_pipeline_names(trigger: AdfTrigger) -> list[str]: return names -def adf_pipeline_to_source_graph(pipeline: AdfPipeline) -> SourceGraph: - """Map a single ADF pipeline to a source-faithful :class:`SourceGraph`.""" +def adf_pipeline_to_source_graph(pipeline: AdfPipeline, definitions: AdfDefinitions | None = None) -> SourceGraph: + """Map a single ADF pipeline to a source-faithful :class:`SourceGraph`. + + ``definitions`` supplies the datasets and linked services the activity mapper + needs to resolve each node's data assets to a physical identity. It is optional + so a caller mapping a pipeline in isolation still works; without it, dataset + references simply stay unresolved (identity ``None``) and fall back to their + neutral signature. It does **not** attach the graph-level lineage block -- + :func:`adf_definitions_to_source_graphs` owns that, once every graph is built. + """ + resolved_definitions = definitions if definitions is not None else AdfDefinitions(pipelines=[]) properties: dict[str, str] = {} if pipeline.folder: properties["folder"] = pipeline.folder @@ -194,17 +221,21 @@ def adf_pipeline_to_source_graph(pipeline: AdfPipeline) -> SourceGraph: parameters=_parameters_to_specs(pipeline.parameters), variables=_variables_to_specs(pipeline.variables), tags=list(pipeline.annotations) if pipeline.annotations else [], - tasks=[_activity_to_node(activity) for activity in pipeline.activities], + tasks=[_activity_to_node(activity, resolved_definitions) for activity in pipeline.activities], properties=properties, raw=pipeline.raw, ) -def _activity_to_node(activity: AdfActivity) -> SourceNode: +def _activity_to_node(activity: AdfActivity, definitions: AdfDefinitions) -> SourceNode: """Map one ADF activity to a discovery node (1:1, source-faithful). Fields are passed explicitly to each node class (rather than unpacking a - shared dict) so the mapping stays type-checked end to end. + shared dict) so the mapping stays type-checked end to end. Data reads / writes + are resolved from the activity's dataset references, and an ``ExecutePipeline`` + records the child pipeline it invokes under the neutral control-edge marker; + both apply to nested nodes too because the branch children below route back + through this function. """ strategy = classify_activity(activity.type) concept = _CONCEPT_BY_TYPE.get(activity.type, CONCEPT_GAP) @@ -214,8 +245,10 @@ def _activity_to_node(activity: AdfActivity) -> SourceNode: ] policy = _policy_to_spec(activity) properties: dict[str, Any] = {STRATEGY_PROPERTY: strategy.value} + _record_invocation(activity, properties) + data_reads, data_writes = activity_data_assets(activity, definitions) - branches = _control_flow_branches(activity) + branches = _control_flow_branches(activity, definitions) if branches is not None: return ContainerNode( source_id=activity.name, @@ -226,6 +259,8 @@ def _activity_to_node(activity: AdfActivity) -> SourceNode: native_type=activity.type, dependencies=dependencies, policy=policy, + data_reads=data_reads, + data_writes=data_writes, properties=properties, raw=activity.raw, branches=branches, @@ -239,6 +274,8 @@ def _activity_to_node(activity: AdfActivity) -> SourceNode: native_type=activity.type, dependencies=dependencies, policy=policy, + data_reads=data_reads, + data_writes=data_writes, properties=properties, raw=activity.raw, reason=f"unsupported ADF activity type {activity.type!r}", @@ -252,12 +289,36 @@ def _activity_to_node(activity: AdfActivity) -> SourceNode: native_type=activity.type, dependencies=dependencies, policy=policy, + data_reads=data_reads, + data_writes=data_writes, properties=properties, raw=activity.raw, ) -def _control_flow_branches(activity: AdfActivity) -> dict[str, list[SourceNode]] | None: +def _record_invocation(activity: AdfActivity, properties: dict[str, Any]) -> None: + """Stash the child pipeline an ``ExecutePipeline`` invokes under neutral keys. + + The callee reference name and ``waitOnCompletion`` flag are read from the same + ADF ``typeProperties`` the convert-time translator reads, then recorded under + the source-neutral :data:`INVOKES_WORKFLOW_PROPERTY` / :data:`INVOKES_WAIT_PROPERTY` + keys so the source-agnostic control-edge derivation can find them without + knowing anything about ADF. An empty / missing callee still records the marker + (with an empty target) so the edge is reported as unresolved rather than dropped. + """ + if activity.type != "ExecutePipeline": + return + type_properties = activity.type_properties or {} + reference = type_properties.get("pipeline", {}) + if isinstance(reference, dict): + callee = reference.get("referenceName", "") or "" + else: + callee = str(reference) + properties[INVOKES_WORKFLOW_PROPERTY] = callee + properties[INVOKES_WAIT_PROPERTY] = bool(type_properties.get("waitOnCompletion", True)) + + +def _control_flow_branches(activity: AdfActivity, definitions: AdfDefinitions) -> dict[str, list[SourceNode]] | None: """Return the labelled branches of a control-flow activity, or ``None`` if it is a leaf. Keyed on the activity **type**, not on whether children happen to be present: @@ -273,16 +334,18 @@ def _control_flow_branches(activity: AdfActivity) -> dict[str, list[SourceNode]] """ if activity.type == "IfCondition": return { - "true": [_activity_to_node(child) for child in (activity.if_true_activities or [])], - "false": [_activity_to_node(child) for child in (activity.if_false_activities or [])], + "true": [_activity_to_node(child, definitions) for child in (activity.if_true_activities or [])], + "false": [_activity_to_node(child, definitions) for child in (activity.if_false_activities or [])], } if activity.type in ("ForEach", "Until"): - return {"body": [_activity_to_node(child) for child in (activity.activities or [])]} + return {"body": [_activity_to_node(child, definitions) for child in (activity.activities or [])]} if activity.type == "Switch": branches: dict[str, list[SourceNode]] = {} for case_value, case_activities in (activity.switch_cases or {}).items(): - branches[case_value] = [_activity_to_node(child) for child in case_activities] - branches["default"] = [_activity_to_node(child) for child in (activity.switch_default_activities or [])] + branches[case_value] = [_activity_to_node(child, definitions) for child in case_activities] + branches["default"] = [ + _activity_to_node(child, definitions) for child in (activity.switch_default_activities or []) + ] return branches return None diff --git a/tests/unit/test_adf_dataset_lineage.py b/tests/unit/test_adf_dataset_lineage.py new file mode 100644 index 0000000..8d3f9a5 --- /dev/null +++ b/tests/unit/test_adf_dataset_lineage.py @@ -0,0 +1,239 @@ +"""Tests for the ADF dataset-lineage resolver (:mod:`flowx.sources.adf.dataset_lineage`). + +Proves the two-tier identity / signature model re-homed from the closed #36 +attempt: a table or a concrete path resolves to a physical ``identity``; a +parameterised reference does not (it never guesses) and falls back to the +structural path signature or the dataset name; and every input / output slot is +captured, not just index 0. +""" + +from __future__ import annotations + +from flowx.models.adf_ast import ( + AdfActivity, + AdfDataset, + AdfDatasetReference, + AdfDefinitions, + AdfLinkedService, +) +from flowx.sources.adf.dataset_lineage import activity_data_assets, resolve_dataset_identity + + +def _definitions(**datasets: AdfDataset) -> AdfDefinitions: + return AdfDefinitions(pipelines=[], datasets=dict(datasets)) + + +def _table_dataset(name: str, *, schema: str, table: str) -> AdfDataset: + return AdfDataset( + name=name, + type="AzureSqlTable", + properties={"typeProperties": {"schema": schema, "table": table}}, + ) + + +def _adls_dataset(name: str, *, file_system: str, folder_path: str, linked_service: str) -> AdfDataset: + return AdfDataset( + name=name, + type="DelimitedText", + properties={ + "typeProperties": {"location": {"fileSystem": file_system, "folderPath": folder_path}}, + "linkedServiceName": {"referenceName": linked_service}, + }, + ) + + +def _adls_linked_service(name: str, *, account: str) -> AdfLinkedService: + return AdfLinkedService( + name=name, + type="AzureBlobFS", + properties={"typeProperties": {"url": f"https://{account}.dfs.core.windows.net"}}, + ) + + +# --------------------------------------------------------------------------- # +# Identity tier +# --------------------------------------------------------------------------- # + + +def test_identity_resolves_schema_and_table() -> None: + definitions = _definitions(ds_orders=_table_dataset("ds_orders", schema="curated", table="orders")) + identity = resolve_dataset_identity(AdfDatasetReference(reference_name="ds_orders"), definitions) + assert identity == "curated.orders" + + +def test_identity_resolves_storage_path_from_linked_service() -> None: + definitions = AdfDefinitions( + pipelines=[], + datasets={ + "ds_raw": _adls_dataset("ds_raw", file_system="data", folder_path="raw/customers", linked_service="ls") + }, + linked_services={"ls": _adls_linked_service("ls", account="contosolake")}, + ) + identity = resolve_dataset_identity(AdfDatasetReference(reference_name="ds_raw"), definitions) + assert identity == "abfss://data@contosolake.dfs.core.windows.net/raw/customers" + + +def test_identity_resolves_dataset_param_from_call_site_literal() -> None: + """A ``@dataset().table`` expression resolves when the call site passes a literal.""" + dataset = AdfDataset( + name="ds_param", + type="AzureSqlTable", + properties={ + "typeProperties": {"schema": "dbo", "table": "@dataset().tbl"}, + "parameters": {"tbl": {"type": "String"}}, + }, + ) + definitions = _definitions(ds_param=dataset) + reference = AdfDatasetReference(reference_name="ds_param", parameters={"tbl": "shipments"}) + assert resolve_dataset_identity(reference, definitions) == "dbo.shipments" + + +def test_parameterised_table_is_not_guessed() -> None: + """An unresolved ``@pipeline()`` expression yields no identity (never a guess, per #36).""" + dataset = AdfDataset( + name="ds_dyn", + type="AzureSqlTable", + properties={"typeProperties": {"schema": "dbo", "table": "@pipeline().parameters.tableName"}}, + ) + definitions = _definitions(ds_dyn=dataset) + assert resolve_dataset_identity(AdfDatasetReference(reference_name="ds_dyn"), definitions) is None + + +def test_unknown_dataset_reference_has_no_identity() -> None: + assert resolve_dataset_identity(AdfDatasetReference(reference_name="missing"), _definitions()) is None + + +# --------------------------------------------------------------------------- # +# activity_data_assets: reads/writes + the two-tier signature +# --------------------------------------------------------------------------- # + + +def test_copy_captures_source_read_and_sink_write_with_identity() -> None: + definitions = _definitions( + ds_src=_table_dataset("ds_src", schema="raw", table="orders"), + ds_dst=_table_dataset("ds_dst", schema="curated", table="orders"), + ) + activity = AdfActivity( + name="Copy Orders", + type="Copy", + type_properties={ + "source": {"referenceName": "ds_src"}, + "sink": {"referenceName": "ds_dst"}, + }, + ) + reads, writes = activity_data_assets(activity, definitions) + assert [(asset.identity, asset.signature) for asset in reads] == [("raw.orders", "raw.orders")] + assert [(asset.identity, asset.signature) for asset in writes] == [("curated.orders", "curated.orders")] + assert reads[0].asset_type == "table" + + +def test_captures_all_inputs_and_outputs_not_just_index_zero() -> None: + """Every ``inputs`` / ``outputs`` slot is captured, not only the first.""" + definitions = _definitions( + ds_in_a=_table_dataset("ds_in_a", schema="raw", table="a"), + ds_in_b=_table_dataset("ds_in_b", schema="raw", table="b"), + ds_out_a=_table_dataset("ds_out_a", schema="curated", table="a"), + ds_out_b=_table_dataset("ds_out_b", schema="curated", table="b"), + ) + activity = AdfActivity( + name="Multi IO", + type="Copy", + inputs=[AdfDatasetReference(reference_name="ds_in_a"), AdfDatasetReference(reference_name="ds_in_b")], + outputs=[AdfDatasetReference(reference_name="ds_out_a"), AdfDatasetReference(reference_name="ds_out_b")], + ) + reads, writes = activity_data_assets(activity, definitions) + assert sorted(asset.identity for asset in reads) == ["raw.a", "raw.b"] + assert sorted(asset.identity for asset in writes) == ["curated.a", "curated.b"] + + +def test_dataset_named_in_both_slot_and_typeproperties_counted_once() -> None: + definitions = _definitions(ds_src=_table_dataset("ds_src", schema="raw", table="orders")) + activity = AdfActivity( + name="Lookup", + type="Lookup", + inputs=[AdfDatasetReference(reference_name="ds_src")], + type_properties={"dataset": {"referenceName": "ds_src"}}, + ) + reads, _ = activity_data_assets(activity, definitions) + assert [asset.identity for asset in reads] == ["raw.orders"] + + +def test_unresolved_reference_falls_back_to_path_signature() -> None: + """A parameterised file reference with a literal folder anchor gets a structural signature.""" + dataset = AdfDataset( + name="ds_wm", + type="DelimitedText", + properties={"typeProperties": {"location": {"fileSystem": "@dataset().fs"}}}, + ) + definitions = _definitions(ds_wm=dataset) + reference = AdfDatasetReference( + reference_name="ds_wm", + parameters={"folderPath": "@concat('watermark/', item().entity)", "fileName": "version.txt"}, + ) + activity = AdfActivity(name="Read WM", type="Lookup", type_properties={"dataset": _ref_dict(reference)}) + reads, _ = activity_data_assets(activity, definitions) + assert reads[0].identity is None + assert reads[0].signature == "FP[watermark|slots=1]/FN[version.txt|slots=0]" + + +def test_unresolved_reference_without_path_anchor_has_empty_signature() -> None: + """No identity and no path anchor -> empty signature, so it never joins (#36 name rule).""" + dataset = AdfDataset( + name="ds_generic", + type="AzureSqlTable", + properties={"typeProperties": {"table": "@pipeline().parameters.t"}}, + ) + definitions = _definitions(ds_generic=dataset) + activity = AdfActivity( + name="Copy", + type="Copy", + type_properties={"source": {"referenceName": "ds_generic"}}, + ) + reads, _ = activity_data_assets(activity, definitions) + assert reads[0].identity is None + # Never the bare dataset name: an empty signature cannot produce a signature-tier join. + assert reads[0].signature == "" + + +def test_distinct_parameter_bindings_of_same_dataset_are_all_retained() -> None: + """Same dataset ref, different params -> distinct physical assets, both kept (not collapsed).""" + dataset = AdfDataset( + name="ds", + type="AzureSqlTable", + properties={ + "typeProperties": {"schema": "raw", "table": "@dataset().tbl"}, + "parameters": {"tbl": {"type": "String"}}, + }, + ) + definitions = _definitions(ds=dataset) + activity = AdfActivity( + name="Multi Bind", + type="Copy", + inputs=[ + AdfDatasetReference(reference_name="ds", parameters={"tbl": "orders"}), + AdfDatasetReference(reference_name="ds", parameters={"tbl": "customers"}), + ], + ) + reads, _ = activity_data_assets(activity, definitions) + assert sorted(asset.identity for asset in reads) == ["raw.customers", "raw.orders"] + + +def test_same_ref_same_params_in_slot_and_typeproperties_still_collapses() -> None: + """The de-dup only fires on an identical binding: one asset, not two.""" + definitions = _definitions(ds_src=_table_dataset("ds_src", schema="raw", table="orders")) + activity = AdfActivity( + name="Lookup", + type="Lookup", + inputs=[AdfDatasetReference(reference_name="ds_src")], + type_properties={"dataset": {"referenceName": "ds_src"}}, + ) + reads, _ = activity_data_assets(activity, definitions) + assert [asset.identity for asset in reads] == ["raw.orders"] + + +def _ref_dict(reference: AdfDatasetReference) -> dict: + """Render a dataset reference the way ADF ``typeProperties`` nests it.""" + payload: dict = {"referenceName": reference.reference_name} + if reference.parameters: + payload["parameters"] = reference.parameters + return payload diff --git a/tests/unit/test_discovery_lineage.py b/tests/unit/test_discovery_lineage.py new file mode 100644 index 0000000..f20f782 --- /dev/null +++ b/tests/unit/test_discovery_lineage.py @@ -0,0 +1,373 @@ +"""Tests for lineage over the shared discovery AST (:mod:`flowx.discovery_lineage`). + +Two layers: + +* the source-neutral derivation itself -- the identity vs signature join tiers, + no self-edges, no duplicates, fan-out, and Switch / ForEach / If recursion -- + driven straight off :class:`SourceGraph` / :class:`SourceNode` values; and +* the end-to-end ADF path -- :func:`adf_definitions_to_source_graphs` populating a + node's reads / writes from the ADF resolver and attaching a derived lineage + block, including nested-in-Switch capture and ExecutePipeline control edges. +""" + +from __future__ import annotations + +from flowx.discovery_lineage import ( + INVOKES_WAIT_PROPERTY, + INVOKES_WORKFLOW_PROPERTY, + build_graph_lineage, + walk_nodes, + with_graph_lineage, +) +from flowx.models.adf_ast import ( + AdfActivity, + AdfDataset, + AdfDefinitions, + AdfPipeline, +) +from flowx.models.discovery import ( + SOURCE_ADF, + ContainerNode, + SourceGraph, + SourceNode, +) +from flowx.models.ir import DataAsset +from flowx.sources.adf.discovery_mapping import adf_definitions_to_source_graphs +from flowx.sources.adf.loader import load_adf_definitions + +# --------------------------------------------------------------------------- # +# Helpers for the source-neutral layer +# --------------------------------------------------------------------------- # + + +def _node(task_key: str, *, reads=None, writes=None, invokes=None, wait=None) -> SourceNode: + properties: dict = {} + if invokes is not None: + properties[INVOKES_WORKFLOW_PROPERTY] = invokes + properties[INVOKES_WAIT_PROPERTY] = wait + return SourceNode( + source_id=task_key, + task_key=task_key, + concept="x", + source=SOURCE_ADF, + data_reads=list(reads or []), + data_writes=list(writes or []), + properties=properties, + ) + + +def _graph(*tasks: SourceNode, name: str = "pl") -> SourceGraph: + return SourceGraph(name=name, source=SOURCE_ADF, tasks=list(tasks)) + + +# --------------------------------------------------------------------------- # +# Data-edge tiers (source-neutral) +# --------------------------------------------------------------------------- # + + +def test_identity_tier_joins_across_different_signatures() -> None: + graph = _graph( + _node("writer", writes=[DataAsset(signature="ds_out", identity="curated.orders")]), + _node("reader", reads=[DataAsset(signature="ds_in_other", identity="curated.orders")]), + ) + edges = build_graph_lineage(graph).data_edges + assert len(edges) == 1 + assert (edges[0].source_task_key, edges[0].target_task_key) == ("writer", "reader") + assert edges[0].match_kind == "identity" + assert edges[0].identity == "curated.orders" + + +def test_signature_tier_when_identity_unresolved() -> None: + graph = _graph( + _node("writer", writes=[DataAsset(signature="FP[wm|slots=1]/FN[v.txt|slots=0]")]), + _node("reader", reads=[DataAsset(signature="FP[wm|slots=1]/FN[v.txt|slots=0]")]), + ) + edges = build_graph_lineage(graph).data_edges + assert len(edges) == 1 + assert edges[0].match_kind == "signature" + assert edges[0].identity is None + + +def test_distinct_identities_do_not_fall_back_to_signature() -> None: + """Two resolved-but-different identities never manufacture a signature edge (#36).""" + graph = _graph( + _node("writer", writes=[DataAsset(signature="shared", identity="a.first")]), + _node("reader", reads=[DataAsset(signature="shared", identity="b.second")]), + ) + assert build_graph_lineage(graph).data_edges == [] + + +def test_fan_out_one_writer_many_readers() -> None: + graph = _graph( + _node("writer", writes=[DataAsset(signature="ds", identity="x.y")]), + _node("reader_one", reads=[DataAsset(signature="ds", identity="x.y")]), + _node("reader_two", reads=[DataAsset(signature="ds", identity="x.y")]), + ) + edges = build_graph_lineage(graph).data_edges + assert {edge.target_task_key for edge in edges} == {"reader_one", "reader_two"} + + +def test_no_self_edge_and_no_duplicates() -> None: + graph = _graph( + _node( + "both", + writes=[DataAsset(signature="ds", identity="x.y"), DataAsset(signature="ds", identity="x.y")], + reads=[DataAsset(signature="ds", identity="x.y")], + ), + _node("reader", reads=[DataAsset(signature="ds", identity="x.y")]), + ) + edges = build_graph_lineage(graph).data_edges + assert [(edge.source_task_key, edge.target_task_key) for edge in edges] == [("both", "reader")] + + +def test_data_edges_recurse_into_switch_and_foreach_branches() -> None: + """A writer buried in a Switch case hands off to a reader in a ForEach body.""" + writer = _node("writer", writes=[DataAsset(signature="ds", identity="x.y")]) + reader = _node("reader", reads=[DataAsset(signature="ds", identity="x.y")]) + switch = ContainerNode( + source_id="sw", + task_key="sw", + concept="switch", + source=SOURCE_ADF, + branches={"caseA": [writer], "default": []}, + ) + loop = ContainerNode( + source_id="fe", + task_key="fe", + concept="loop", + source=SOURCE_ADF, + branches={"body": [reader]}, + ) + edges = build_graph_lineage(_graph(switch, loop)).data_edges + assert [(edge.source_task_key, edge.target_task_key) for edge in edges] == [("writer", "reader")] + + +def test_walk_nodes_visits_every_branch_including_default() -> None: + switch = ContainerNode( + source_id="sw", + task_key="sw", + concept="switch", + source=SOURCE_ADF, + branches={"caseA": [_node("in_case")], "default": [_node("in_default")]}, + ) + keys = [node.task_key for node in walk_nodes([switch])] + assert keys == ["sw", "in_case", "in_default"] + + +# --------------------------------------------------------------------------- # +# Control edges (source-neutral) +# --------------------------------------------------------------------------- # + + +def test_control_edge_from_invocation_marker() -> None: + graph = _graph(_node("call", invokes="child", wait=True), name="parent") + edges = build_graph_lineage(graph).control_edges + assert len(edges) == 1 + assert (edges[0].source_workflow, edges[0].target_workflow) == ("parent", "child") + assert edges[0].wait_for_completion is True + assert edges[0].resolved is True + + +def test_control_edge_unresolved_callee_recorded_not_dropped() -> None: + graph = _graph(_node("call", invokes="", wait=True), name="parent") + edges = build_graph_lineage(graph).control_edges + assert len(edges) == 1 + assert edges[0].target_workflow == "" + assert edges[0].resolved is False + + +def test_control_edge_self_call_is_dropped() -> None: + graph = _graph(_node("call", invokes="parent", wait=True), name="parent") + assert build_graph_lineage(graph).control_edges == [] + + +def test_with_graph_lineage_returns_new_graph_and_leaves_input_untouched() -> None: + graph = _graph(_node("call", invokes="child", wait=False), name="parent") + result = with_graph_lineage(graph) + assert graph.lineage is None + assert result is not graph + assert len(result.lineage.control_edges) == 1 + + +# --------------------------------------------------------------------------- # +# End-to-end through the ADF mapper +# --------------------------------------------------------------------------- # + + +def _table_dataset(name: str, *, schema: str, table: str) -> AdfDataset: + return AdfDataset( + name=name, + type="AzureSqlTable", + properties={"typeProperties": {"schema": schema, "table": table}}, + ) + + +def test_adf_within_pipeline_identity_handoff() -> None: + """A Copy writes curated.orders; a later Lookup reads it -> one identity data edge.""" + definitions = AdfDefinitions( + pipelines=[ + AdfPipeline( + name="etl", + activities=[ + AdfActivity( + name="Load Orders", + type="Copy", + type_properties={ + "source": {"referenceName": "ds_raw"}, + "sink": {"referenceName": "ds_curated"}, + }, + ), + AdfActivity( + name="Check Orders", + type="Lookup", + type_properties={"dataset": {"referenceName": "ds_curated"}}, + ), + ], + ) + ], + datasets={ + "ds_raw": _table_dataset("ds_raw", schema="raw", table="orders"), + "ds_curated": _table_dataset("ds_curated", schema="curated", table="orders"), + }, + ) + graph = adf_definitions_to_source_graphs(definitions)[0] + assert graph.lineage is not None + data_edges = graph.lineage.data_edges + assert len(data_edges) == 1 + edge = data_edges[0] + assert (edge.source_task_key, edge.target_task_key) == ("Load Orders", "Check Orders") + assert edge.match_kind == "identity" + assert edge.identity == "curated.orders" + + +def test_adf_opaque_references_do_not_produce_a_false_data_edge() -> None: + """Two unresolvable references sharing a name must NOT join on that name (#36).""" + opaque = AdfDataset( + name="ds_opaque", + type="AzureSqlTable", + properties={"typeProperties": {"table": "@pipeline().parameters.t"}}, + ) + definitions = AdfDefinitions( + pipelines=[ + AdfPipeline( + name="etl", + activities=[ + AdfActivity( + name="Write Opaque", + type="Copy", + type_properties={"sink": {"referenceName": "ds_opaque"}}, + ), + AdfActivity( + name="Read Opaque", + type="Lookup", + type_properties={"dataset": {"referenceName": "ds_opaque"}}, + ), + ], + ) + ], + datasets={"ds_opaque": opaque}, + ) + graph = adf_definitions_to_source_graphs(definitions)[0] + assert graph.lineage is not None + # Both references resolve to no identity and no path anchor -> empty signature -> no edge. + assert graph.lineage.data_edges == [] + writer = next(node for node in walk_nodes(graph.tasks) if node.task_key == "Write Opaque") + assert writer.data_writes[0].signature == "" + + +def test_adf_execute_pipeline_produces_control_edge() -> None: + definitions = AdfDefinitions( + pipelines=[ + AdfPipeline( + name="parent", + activities=[ + AdfActivity( + name="Run Child", + type="ExecutePipeline", + type_properties={ + "pipeline": {"referenceName": "child"}, + "waitOnCompletion": False, + }, + ), + ], + ), + AdfPipeline(name="child", activities=[]), + ], + ) + parent = adf_definitions_to_source_graphs(definitions)[0] + assert parent.lineage is not None + control_edges = parent.lineage.control_edges + assert len(control_edges) == 1 + edge = control_edges[0] + assert (edge.source_workflow, edge.target_workflow) == ("parent", "child") + assert edge.via_task_key == "Run Child" + assert edge.wait_for_completion is False + + +def test_adf_switch_nested_copy_reads_and_writes_are_captured() -> None: + """A Copy inside a Switch case carries resolved reads/writes; its hand-off is found.""" + definitions = AdfDefinitions( + pipelines=[ + AdfPipeline( + name="switched", + activities=[ + AdfActivity( + name="Route", + type="Switch", + switch_cases={ + "sql": [ + AdfActivity( + name="Copy From SQL", + type="Copy", + type_properties={ + "source": {"referenceName": "ds_raw"}, + "sink": {"referenceName": "ds_stage"}, + }, + ) + ] + }, + switch_default_activities=[ + AdfActivity( + name="Read Stage", + type="Lookup", + type_properties={"dataset": {"referenceName": "ds_stage"}}, + ) + ], + ), + ], + ) + ], + datasets={ + "ds_raw": _table_dataset("ds_raw", schema="raw", table="events"), + "ds_stage": _table_dataset("ds_stage", schema="stage", table="events"), + }, + ) + graph = adf_definitions_to_source_graphs(definitions)[0] + + # The nested Copy carries its resolved reads/writes. + nested = {node.task_key: node for node in walk_nodes(graph.tasks)} + copy_node = nested["Copy From SQL"] + assert [asset.identity for asset in copy_node.data_reads] == ["raw.events"] + assert [asset.identity for asset in copy_node.data_writes] == ["stage.events"] + + # And the writer (Switch case) -> reader (Switch default) hand-off is derived. + assert graph.lineage is not None + data_edges = graph.lineage.data_edges + assert len(data_edges) == 1 + assert (data_edges[0].source_task_key, data_edges[0].target_task_key) == ("Copy From SQL", "Read Stage") + assert data_edges[0].identity == "stage.events" + + +def test_shared_execute_pipeline_fixture_reproduces_hash36_control_edges(fixtures_dir) -> None: + """The nested-ExecutePipeline fixture yields exactly #36's three caller->callee edges.""" + graphs = adf_definitions_to_source_graphs(load_adf_definitions(fixtures_dir)) + graph = next(g for g in graphs if g.name == "pipeline_execute_pipeline_nested") + assert graph.lineage is not None + edges = { + (edge.source_workflow, edge.target_workflow, edge.wait_for_completion) for edge in graph.lineage.control_edges + } + assert edges == { + ("pipeline_execute_pipeline_nested", "pipeline_copy_sql_to_delta", True), + ("pipeline_execute_pipeline_nested", "pipeline_notebook_with_params", True), + ("pipeline_execute_pipeline_nested", "pipeline_delete_recursive", False), + } diff --git a/tests/unit/test_discovery_model.py b/tests/unit/test_discovery_model.py index 2e1191e..df3cec6 100644 --- a/tests/unit/test_discovery_model.py +++ b/tests/unit/test_discovery_model.py @@ -26,7 +26,7 @@ SourceGraph, SourceNode, ) -from flowx.models.ir import DataAsset +from flowx.models.ir import ControlEdge, DataAsset, DataEdge, Lineage def test_gap_node_defaults_to_gap_concept(): @@ -214,3 +214,38 @@ def test_empty_graph_emits_lists_and_dicts_never_null(): assert serialised["variables"] == {} # Airflow has no graph-scoped variables; the field stays empty rather than absent. assert source_graph_from_dict(serialised).variables == {} + + +def test_lineage_block_round_trips_and_is_absent_when_none(): + """A graph's derived lineage survives serialise<->deserialise; None stays absent.""" + graph = SourceGraph( + name="pl", + source=SOURCE_ADF, + lineage=Lineage( + control_edges=[ + ControlEdge( + source_workflow="pl", + target_workflow="child", + via_task_key="Run Child", + wait_for_completion=False, + ) + ], + data_edges=[ + DataEdge( + source_task_key="writer", + target_task_key="reader", + match_kind="identity", + match_key="curated.orders", + identity="curated.orders", + asset_type="table", + ) + ], + ), + ) + reloaded = source_graph_from_dict(json.loads(json.dumps(source_graph_to_dict(graph)))) + assert reloaded == graph + + # A graph with no derived lineage omits the key entirely and rehydrates to None. + bare = source_graph_to_dict(SourceGraph(name="bare", source=SOURCE_ADF)) + assert "lineage" not in bare + assert source_graph_from_dict(bare).lineage is None