diff --git a/src/flowx/discovery_inventory.py b/src/flowx/discovery_inventory.py new file mode 100644 index 0000000..975e5cc --- /dev/null +++ b/src/flowx/discovery_inventory.py @@ -0,0 +1,156 @@ +"""Source-agnostic projection of the shared discovery AST to ``inventory.json``. + +The discover phase writes ``metadata/inventory.json`` and the reporting layer +(:mod:`flowx.reporting.coverage`) and MCP surface (:mod:`flowx.mcp.runner`) read +it back. Historically each source built that JSON straight from its own AST, so +the shape drifted per source. This module is the single place that turns the +shared discovery AST (:mod:`flowx.models.discovery`) into the inventory shape, so +every source that maps onto :class:`~flowx.models.discovery.SourceGraph` emits the +*same* top-level document -- ``{source, source_dir, pipelines, summary}`` -- from +one code path. + +The projection is deliberately small and additive over the historical ADF shape: + +* top level gains a ``source`` discriminator (``"adf"`` / ``"airflow"``); +* each activity keeps its byte-compatible ``name`` / ``type`` / ``strategy`` (and + ``depends_on`` names when present) and gains the standardised + ``original_type``, ``dependencies`` (upstream **with conditions**), and the + verbatim per-node ``raw``; +* the ``summary`` keeps the historical count block. + +A node's translation ``strategy`` is a Databricks-*target* classification rather +than a source concept, so it is not a typed field on the discovery AST. By +convention a mapper stashes it under ``node.properties["strategy"]`` (see +:data:`STRATEGY_PROPERTY`); this module reads it there. Anything else a source +wants to layer on -- Airflow's audited-count block, findings, reconciliation +status -- rides additively on top of this base and is out of scope here. +""" + +from __future__ import annotations + +from typing import Any + +from flowx.models.discovery import ContainerNode, SourceGraph, SourceNode + +# Well-known property key under which a mapper records a node's Databricks-target +# translation strategy ("deterministic" / "agentic" / "unsupported"). Kept in the +# free-form properties seam because strategy is a target concern, not a shared +# source concept, so it earns no typed field on the discovery AST. +STRATEGY_PROPERTY = "strategy" + +_DETERMINISTIC = "deterministic" +_AGENTIC = "agentic" + + +def build_source_inventory( + graphs: list[SourceGraph], + *, + source: str, + source_dir: str, + include_empty_pipelines: bool = True, +) -> dict[str, Any]: + """Project a list of source graphs into the ``inventory.json`` document. + + Args: + graphs: The source workflows to inventory, already mapped onto the shared + discovery AST. + source: Source discriminator for the top-level ``source`` field + (``SOURCE_ADF`` / ``SOURCE_AIRFLOW``). + source_dir: Original source directory, echoed back for provenance. + include_empty_pipelines: When ``False``, a graph that contributes no + activities is left out of the ``pipelines`` list but still counted in + ``summary.pipeline_count`` -- this reproduces ADF's long-standing + behaviour of omitting zero-activity pipelines from the per-pipeline + listing while still reporting them in the totals. Sources that list + every workflow (Airflow) leave this ``True``. + + Returns: + A JSON-friendly dict with ``source``, ``source_dir``, ``pipelines`` and + ``summary`` keys. + """ + pipeline_entries: list[dict[str, Any]] = [] + deterministic = 0 + agentic = 0 + unsupported = 0 + + for graph in graphs: + flattened = _flatten_nodes(graph.tasks) + for node in flattened: + strategy = node.properties.get(STRATEGY_PROPERTY) + if strategy == _DETERMINISTIC: + deterministic += 1 + elif strategy == _AGENTIC: + agentic += 1 + else: + unsupported += 1 + if flattened or include_empty_pipelines: + pipeline_entries.append( + { + "name": graph.name, + "activities": [_activity_entry(node) for node in flattened], + } + ) + + total = deterministic + agentic + unsupported + coverage_pct = round((deterministic + agentic) / total * 100, 1) if total else 0.0 + + return { + "source": source, + "source_dir": source_dir, + "pipelines": pipeline_entries, + "summary": { + "pipeline_count": len(graphs), + "activity_count": total, + "deterministic_count": deterministic, + "agentic_count": agentic, + "unsupported_count": unsupported, + "coverage_pct": coverage_pct, + }, + } + + +def _flatten_nodes(nodes: list[SourceNode]) -> list[SourceNode]: + """Flatten container branches into one depth-first activity list. + + Order is parent, then each branch's children in the branch's own insertion + order, recursively -- so an ``IfCondition``'s ``true`` branch precedes its + ``false`` branch and a ``Switch``'s cases precede its ``default``, matching the + order the source declared them. + """ + flattened: list[SourceNode] = [] + for node in nodes: + flattened.append(node) + if isinstance(node, ContainerNode): + for children in node.branches.values(): + flattened.extend(_flatten_nodes(children)) + return flattened + + +def _activity_entry(node: SourceNode) -> dict[str, Any]: + """Build one per-activity inventory entry from a node. + + The first three keys (plus ``depends_on`` when the node has dependencies) + reproduce the historical ADF activity shape byte-for-byte; the rest are the + additive standardised fields. + """ + entry: dict[str, Any] = { + "name": node.name if node.name is not None else node.task_key, + "type": node.native_type, + "strategy": node.properties.get(STRATEGY_PROPERTY), + } + upstream_names = [dependency.upstream for dependency in node.dependencies] + if upstream_names: + entry["depends_on"] = upstream_names + + entry["original_type"] = node.native_type + entry["dependencies"] = [ + { + "upstream": dependency.upstream, + "conditions": list(dependency.conditions), + "resolved": dependency.resolved, + } + for dependency in node.dependencies + ] + if node.raw is not None: + entry["raw"] = node.raw + return entry diff --git a/src/flowx/sources/adf/discovery_mapping.py b/src/flowx/sources/adf/discovery_mapping.py new file mode 100644 index 0000000..5841492 --- /dev/null +++ b/src/flowx/sources/adf/discovery_mapping.py @@ -0,0 +1,331 @@ +"""Map the ADF AST onto the shared, source-faithful discovery AST. + +This is the ADF half of the discovery contract (issue #62): it turns the typed +ADF AST (:mod:`flowx.models.adf_ast`) into the shared +:class:`~flowx.models.discovery.SourceGraph` model that both ADF and Airflow +align to. The mapping is deliberately **1:1 and lossless**: + +* every ADF activity becomes exactly one discovery node -- no motif collapse, + no merging; +* the activity's own type string is kept verbatim as ``native_type``; +* the verbatim source dict is preserved on every node (``raw``) and on the graph; +* dependency edges keep **all** of their outcome conditions (a ``dependsOn`` with + ``["Succeeded", "Skipped"]`` stays a two-condition edge); +* anything without a shared typed home -- the ADF folder, retry policy detail, + 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. + +The control-flow container shape follows :class:`~flowx.models.discovery.ContainerNode`: +an ``IfCondition`` becomes ``{"true": [...], "false": [...]}``, a ``ForEach`` / +``Until`` becomes ``{"body": [...]}``, and a ``Switch`` becomes +``{"": [...], "default": [...]}``. Branch insertion order matches the +order ADF declares the branches so a downstream flatten reproduces source order. +""" + +from __future__ import annotations + +from typing import Any + +from flowx.discovery_inventory import STRATEGY_PROPERTY +from flowx.models.adf_ast import AdfActivity, AdfDefinitions, AdfParameter, AdfPipeline, AdfTrigger, AdfVariable +from flowx.models.discovery import ( + CONCEPT_BRANCH, + CONCEPT_COPY_DATA, + CONCEPT_GAP, + CONCEPT_LOOP, + CONCEPT_NOTEBOOK, + CONCEPT_QUERY, + CONCEPT_RUN_WORKFLOW, + CONCEPT_SCRIPT, + CONCEPT_SET_VARIABLE, + CONCEPT_SWITCH, + CONCEPT_WAIT, + SOURCE_ADF, + ContainerNode, + GapNode, + ParameterSpec, + PolicySpec, + ScheduleSpec, + SourceDependency, + SourceGraph, + SourceNode, +) +from flowx.sources.adf.loader import classify_activity + +# Neutral trigger category for each ADF trigger type. Anything unrecognised maps +# to ``""`` (unknown) with the ADF type preserved verbatim in the schedule +# extensions, so no trigger information is lost even for a type not listed here. +_TRIGGER_KIND: dict[str, str] = { + "ScheduleTrigger": "schedule", + "TumblingWindowTrigger": "interval", + "BlobEventsTrigger": "file_arrival", + "CustomEventsTrigger": "event", +} + +# Neutral concept for each ADF activity type. The discovery concept vocabulary is +# intentionally non-exhaustive: ADF types with no shared concept (WebActivity, +# Delete, Filter, and the agentic-only types) map to CONCEPT_GAP, which here means +# "no shared concept applies" -- independent of the translation strategy, which is +# recorded separately under properties[STRATEGY_PROPERTY]. +_CONCEPT_BY_TYPE: dict[str, str] = { + "Copy": CONCEPT_COPY_DATA, + "DatabricksNotebook": CONCEPT_NOTEBOOK, + "DatabricksSparkJar": CONCEPT_SCRIPT, + "DatabricksSparkPython": CONCEPT_SCRIPT, + "DatabricksJob": CONCEPT_RUN_WORKFLOW, + "ExecutePipeline": CONCEPT_RUN_WORKFLOW, + "ForEach": CONCEPT_LOOP, + "Until": CONCEPT_LOOP, + "IfCondition": CONCEPT_BRANCH, + "Switch": CONCEPT_SWITCH, + "SetVariable": CONCEPT_SET_VARIABLE, + "AppendVariable": CONCEPT_SET_VARIABLE, + "Wait": CONCEPT_WAIT, + "Lookup": CONCEPT_QUERY, +} + + +def adf_definitions_to_source_graphs(definitions: AdfDefinitions) -> list[SourceGraph]: + """Map every pipeline in *definitions* to a shared :class:`SourceGraph`. + + 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. + """ + graphs = [adf_pipeline_to_source_graph(pipeline) for pipeline in definitions.pipelines] + _attach_schedules(graphs, definitions.triggers) + return graphs + + +def _attach_schedules(graphs: list[SourceGraph], triggers: list[AdfTrigger]) -> None: + """Populate each graph's ``schedule`` from the triggers that reference it. + + A trigger can drive several pipelines and a pipeline can be driven by several + triggers. The first trigger to reference a pipeline becomes its typed + :attr:`SourceGraph.schedule`; any further triggers for the same pipeline are + preserved verbatim under ``schedule.extensions["additional_triggers"]`` so + nothing is lost. Pipeline references are matched case-insensitively, mirroring + ADF's case-insensitive identifier semantics. + + A **fresh** :class:`ScheduleSpec` is built for each pipeline assignment rather + than sharing one instance across every pipeline a trigger references: a shared + instance would let a later trigger's mutation (appending to + ``additional_triggers``) leak onto every other pipeline that trigger touched. + """ + graphs_by_name = {graph.name: graph for graph in graphs} + graphs_by_lower = {graph.name.lower(): graph for graph in graphs} + + for trigger in triggers: + for pipeline_name in _trigger_pipeline_names(trigger): + graph = graphs_by_name.get(pipeline_name) or graphs_by_lower.get(pipeline_name.lower()) + if graph is None: + continue + if graph.schedule is None: + # Per-pipeline instance: no shared mutable state across pipelines. + graph.schedule = _trigger_to_schedule(trigger) + else: + additional = graph.schedule.extensions.setdefault("additional_triggers", []) + additional.append( + { + "trigger_name": trigger.name, + "trigger_type": trigger.type, + "properties": trigger.properties, + } + ) + + +def _trigger_to_schedule(trigger: AdfTrigger) -> ScheduleSpec: + """Map a single ADF trigger to a source-faithful :class:`ScheduleSpec`. + + The recurrence payload is kept as-given: schedule triggers nest it under + ``typeProperties.recurrence`` while tumbling-window / event triggers put their + detail directly in ``typeProperties``, so whichever is present becomes the + verbatim ``expression``. The full trigger ``properties`` block also rides in + ``extensions`` so the mapping is lossless even for trigger detail with no typed + home yet (pipeline parameters, runtime state, annotations). + """ + type_properties = trigger.properties.get("typeProperties") or {} + recurrence = type_properties.get("recurrence") if isinstance(type_properties, dict) else None + if isinstance(recurrence, dict): + expression: Any = recurrence + timezone = recurrence.get("timeZone") + else: + expression = type_properties or None + timezone = type_properties.get("timeZone") if isinstance(type_properties, dict) else None + + return ScheduleSpec( + kind=_TRIGGER_KIND.get(trigger.type, ""), + expression=expression, + timezone=timezone, + extensions={ + "trigger_name": trigger.name, + "trigger_type": trigger.type, + "properties": trigger.properties, + }, + ) + + +def _trigger_pipeline_names(trigger: AdfTrigger) -> list[str]: + """Collect the names of the pipelines a trigger references, in order.""" + names: list[str] = [] + for reference in trigger.pipelines or []: + pipeline_reference = reference.get("pipelineReference") if isinstance(reference, dict) else None + name = pipeline_reference.get("referenceName") if isinstance(pipeline_reference, dict) else None + if name: + names.append(name) + return names + + +def adf_pipeline_to_source_graph(pipeline: AdfPipeline) -> SourceGraph: + """Map a single ADF pipeline to a source-faithful :class:`SourceGraph`.""" + properties: dict[str, str] = {} + if pipeline.folder: + properties["folder"] = pipeline.folder + + return SourceGraph( + name=pipeline.name, + source=SOURCE_ADF, + 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], + properties=properties, + raw=pipeline.raw, + ) + + +def _activity_to_node(activity: AdfActivity) -> 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. + """ + strategy = classify_activity(activity.type) + concept = _CONCEPT_BY_TYPE.get(activity.type, CONCEPT_GAP) + dependencies = [ + SourceDependency(upstream=dependency.activity, conditions=list(dependency.dependency_conditions)) + for dependency in (activity.depends_on or []) + ] + policy = _policy_to_spec(activity) + properties: dict[str, Any] = {STRATEGY_PROPERTY: strategy.value} + + branches = _control_flow_branches(activity) + if branches is not None: + return ContainerNode( + source_id=activity.name, + task_key=activity.name, + concept=concept, + source=SOURCE_ADF, + name=activity.name, + native_type=activity.type, + dependencies=dependencies, + policy=policy, + properties=properties, + raw=activity.raw, + branches=branches, + ) + if strategy.value == "unsupported": + return GapNode( + source_id=activity.name, + task_key=activity.name, + source=SOURCE_ADF, + name=activity.name, + native_type=activity.type, + dependencies=dependencies, + policy=policy, + properties=properties, + raw=activity.raw, + reason=f"unsupported ADF activity type {activity.type!r}", + ) + return SourceNode( + source_id=activity.name, + task_key=activity.name, + concept=concept, + source=SOURCE_ADF, + name=activity.name, + native_type=activity.type, + dependencies=dependencies, + policy=policy, + properties=properties, + raw=activity.raw, + ) + + +def _control_flow_branches(activity: AdfActivity) -> 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: + a control-flow activity always maps to a :class:`ContainerNode`, and every + branch it declares is always represented -- an empty branch is present-but-empty, + never omitted. So an empty ``IfCondition`` still yields ``{"true": [], "false": []}`` + and a one-sided ``If`` keeps its empty ``false`` branch, rather than collapsing to + a plain node and losing the control-flow structure. Returning ``None`` (not an + empty dict) is what tells the caller the activity is a leaf. + + Branch order matches ADF's own declaration order (true before false, cases + before default) so a downstream flatten walks children in source order. + """ + 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 [])], + } + if activity.type in ("ForEach", "Until"): + return {"body": [_activity_to_node(child) 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 [])] + return branches + return None + + +def _policy_to_spec(activity: AdfActivity) -> PolicySpec | None: + """Map an ADF retry/timeout policy to a :class:`PolicySpec`, if present. + + Retry count and interval are already numeric in ADF and map to typed fields. + The ADF timeout is an ISO-8601 duration *string* -- normalising it to seconds + is a target concern, so it rides verbatim in ``extensions`` alongside the + ``secureInput`` / ``secureOutput`` flags rather than being guessed at here. + """ + policy = activity.policy + if policy is None: + return None + extensions: dict[str, object] = {} + if policy.timeout is not None: + extensions["timeout"] = policy.timeout + if policy.secure_input: + extensions["secure_input"] = True + if policy.secure_output: + extensions["secure_output"] = True + return PolicySpec( + max_retries=policy.retry, + retry_interval_seconds=policy.retry_interval_in_seconds, + extensions=extensions, + ) + + +def _parameters_to_specs(parameters: dict[str, AdfParameter] | None) -> dict[str, ParameterSpec]: + """Map ADF pipeline parameters to shared parameter specs.""" + if not parameters: + return {} + return { + name: ParameterSpec(type=parameter.type, default=parameter.default_value) + for name, parameter in parameters.items() + } + + +def _variables_to_specs(variables: dict[str, AdfVariable] | None) -> dict[str, ParameterSpec]: + """Map ADF pipeline variables to shared parameter specs (same shape).""" + if not variables: + return {} + return { + name: ParameterSpec(type=variable.type, default=variable.default_value) for name, variable in variables.items() + } diff --git a/src/flowx/sources/adf/loader.py b/src/flowx/sources/adf/loader.py index 5fcc5d6..6499071 100644 --- a/src/flowx/sources/adf/loader.py +++ b/src/flowx/sources/adf/loader.py @@ -8,7 +8,6 @@ import logging import re import shutil -from datetime import datetime, timezone from pathlib import Path from typing import Any @@ -789,50 +788,6 @@ def _classify_activities( _classify_activities(pipeline_name, switch_children, items) -# --------------------------------------------------------------------------- -# Serialisation helpers -# --------------------------------------------------------------------------- - - -def _inventory_to_dict(inventory: Inventory, source_dir: str) -> dict[str, Any]: - """Serialise an :class:`Inventory` to a JSON-friendly dictionary. - - Args: - inventory: The inventory to serialise. - source_dir: Original source directory path (for provenance). - - Returns: - Dictionary suitable for ``json.dumps``. - """ - pipeline_map: dict[str, list[dict[str, Any]]] = {} - for item in inventory.items: - entry: dict[str, Any] = { - "name": item.activity_name, - "type": item.activity_type, - "strategy": item.strategy.value, - } - if item.depends_on: - entry["depends_on"] = item.depends_on - pipeline_map.setdefault(item.pipeline_name, []).append(entry) - - total = inventory.deterministic_count + inventory.agentic_count + inventory.unsupported_count - coverage_pct = round((inventory.deterministic_count + inventory.agentic_count) / total * 100, 1) if total else 0.0 - - return { - "source_dir": source_dir, - "generated_at": datetime.now(timezone.utc).isoformat(), - "pipelines": [{"name": pname, "activities": acts} for pname, acts in pipeline_map.items()], - "summary": { - "pipeline_count": inventory.pipeline_count, - "activity_count": total, - "deterministic_count": inventory.deterministic_count, - "agentic_count": inventory.agentic_count, - "unsupported_count": inventory.unsupported_count, - "coverage_pct": coverage_pct, - }, - } - - # --------------------------------------------------------------------------- # Profile complexity report (CSV) # --------------------------------------------------------------------------- @@ -1095,7 +1050,12 @@ def main(argv: list[str] | None = None) -> int: ) logger.info("Filtered to pipeline: %s", args.pipeline) - inventory = build_inventory(definitions) + # Inventory JSON is projected from the shared discovery AST via the + # source-agnostic emitter (so ADF and Airflow emit one shape); imported here + # to avoid a module-level cycle (discovery_mapping imports this loader). + from flowx.discovery_inventory import build_source_inventory + from flowx.models.discovery import SOURCE_ADF + from flowx.sources.adf.discovery_mapping import adf_definitions_to_source_graphs output_dir: Path = args.output_dir.resolve() clear_stale_outputs(output_dir) @@ -1103,7 +1063,15 @@ def main(argv: list[str] | None = None) -> int: metadata_dir.mkdir(parents=True, exist_ok=True) inventory_path = metadata_dir / "inventory.json" - inventory_dict = _inventory_to_dict(inventory, str(args.source_dir)) + source_graphs = adf_definitions_to_source_graphs(definitions) + inventory_dict = build_source_inventory( + source_graphs, + source=SOURCE_ADF, + source_dir=str(args.source_dir), + # ADF has historically omitted zero-activity pipelines from the per-pipeline + # listing while still counting them in summary.pipeline_count; preserve that. + include_empty_pipelines=False, + ) inventory_path.write_text(json.dumps(inventory_dict, indent=2), encoding="utf-8") logger.info("Wrote inventory to %s", inventory_path) diff --git a/tests/resources/golden/adf_fixture_coverage.json b/tests/resources/golden/adf_fixture_coverage.json new file mode 100644 index 0000000..c36912d --- /dev/null +++ b/tests/resources/golden/adf_fixture_coverage.json @@ -0,0 +1,1162 @@ +[ + { + "activities": 13, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 13, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 31, + "complexity_size": "XL", + "control_flow_activities": 4, + "coverage_pct": 100.0, + "databricks_native_activities": 4, + "datasets": 3, + "deterministic_activities": 13, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 5, + "pipeline": "pipeline_all_activity_types", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 8, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 8, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 1, + "complexity_score": 23, + "complexity_size": "L", + "control_flow_activities": 3, + "coverage_pct": 100.0, + "databricks_native_activities": 1, + "datasets": 2, + "deterministic_activities": 8, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 4, + "pipeline": "pipeline_all_dependency_conditions", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 5, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 5, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 1, + "complexity_score": 14, + "complexity_size": "M", + "control_flow_activities": 3, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 1, + "deterministic_activities": 5, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 2, + "pipeline": "pipeline_append_variable_loop", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 8, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 8, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 3, + "complexity_score": 26, + "complexity_size": "L", + "control_flow_activities": 1, + "coverage_pct": 100.0, + "databricks_native_activities": 3, + "datasets": 4, + "deterministic_activities": 8, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 2, + "migration_status": "included", + "other_activities": 4, + "pipeline": "pipeline_complex_etl", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 15, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 15, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 1, + "complexity_score": 39, + "complexity_size": "XL", + "control_flow_activities": 7, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 0, + "deterministic_activities": 15, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 8, + "pipeline": "pipeline_complex_orchestration", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 1, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 1, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 7, + "complexity_size": "M", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 2, + "deterministic_activities": 1, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 2, + "migration_status": "included", + "other_activities": 1, + "pipeline": "pipeline_copy_csv_to_delta", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 1, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 1, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 6, + "complexity_size": "M", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 2, + "deterministic_activities": 1, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 1, + "pipeline": "pipeline_copy_parquet_to_delta", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 1, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 1, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 6, + "complexity_size": "M", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 2, + "deterministic_activities": 1, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 1, + "pipeline": "pipeline_copy_sql_to_delta", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 2, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 2, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 8, + "complexity_size": "M", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 2, + "deterministic_activities": 2, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 2, + "pipeline": "pipeline_delete_recursive", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 3, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 3, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 9, + "complexity_size": "M", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 0, + "deterministic_activities": 3, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 3, + "pipeline": "pipeline_execute_pipeline_nested", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 4, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 4, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 10, + "complexity_size": "M", + "control_flow_activities": 2, + "coverage_pct": 100.0, + "databricks_native_activities": 1, + "datasets": 1, + "deterministic_activities": 4, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 1, + "pipeline": "pipeline_filter_array", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 7, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 7, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 18, + "complexity_size": "L", + "control_flow_activities": 3, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 0, + "deterministic_activities": 7, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 4, + "pipeline": "pipeline_foreach_switch", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 2, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 2, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 7, + "complexity_size": "M", + "control_flow_activities": 1, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 2, + "deterministic_activities": 2, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 1, + "pipeline": "pipeline_foreach_with_copy", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 5, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 5, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 15, + "complexity_size": "M", + "control_flow_activities": 1, + "coverage_pct": 100.0, + "databricks_native_activities": 1, + "datasets": 2, + "deterministic_activities": 5, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 3, + "pipeline": "pipeline_if_condition_branching", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 3, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 3, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 1, + "complexity_score": 12, + "complexity_size": "M", + "control_flow_activities": 1, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 3, + "deterministic_activities": 3, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 2, + "pipeline": "pipeline_lookup_and_foreach", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 6, + "agentic_activities": 4, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 6, + "code_attached_coverage_pct": 83.3, + "collapsible_patterns": 0, + "complexity_score": 19, + "complexity_size": "L", + "control_flow_activities": 0, + "coverage_pct": 83.3, + "databricks_native_activities": 0, + "datasets": 1, + "deterministic_activities": 1, + "deterministic_coverage_pct": 16.7, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 6, + "pipeline": "pipeline_mixed_agentic", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 4, + "unresolved_agentic_count": 0, + "unsupported_activities": 1 + }, + { + "activities": 6, + "agentic_activities": 1, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 6, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 1, + "complexity_score": 20, + "complexity_size": "L", + "control_flow_activities": 1, + "coverage_pct": 100.0, + "databricks_native_activities": 1, + "datasets": 3, + "deterministic_activities": 5, + "deterministic_coverage_pct": 83.3, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 4, + "pipeline": "pipeline_mixed_deterministic_agentic", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 1, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 1, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 1, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 2, + "complexity_size": "S", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 1, + "datasets": 0, + "deterministic_activities": 1, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 0, + "pipeline": "pipeline_notebook_basic", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 1, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 1, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 2, + "complexity_size": "S", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 1, + "datasets": 0, + "deterministic_activities": 1, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 0, + "pipeline": "pipeline_notebook_with_params", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 5, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 5, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 10, + "complexity_size": "M", + "control_flow_activities": 4, + "coverage_pct": 100.0, + "databricks_native_activities": 1, + "datasets": 0, + "deterministic_activities": 5, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 0, + "pipeline": "pipeline_set_variable_chain", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 1, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 1, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 2, + "complexity_size": "S", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 1, + "datasets": 0, + "deterministic_activities": 1, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 0, + "pipeline": "pipeline_spark_jar_job", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 2, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 2, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 4, + "complexity_size": "S", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 2, + "datasets": 0, + "deterministic_activities": 2, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 2, + "migration_status": "included", + "other_activities": 0, + "pipeline": "pipeline_spark_python_job", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 7, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 7, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 23, + "complexity_size": "L", + "control_flow_activities": 2, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 4, + "deterministic_activities": 7, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 5, + "pipeline": "pipeline_switch_multi_case", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 5, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 5, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 12, + "complexity_size": "M", + "control_flow_activities": 2, + "coverage_pct": 100.0, + "databricks_native_activities": 1, + "datasets": 0, + "deterministic_activities": 5, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 2, + "pipeline": "pipeline_wait_between_steps", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 3, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 3, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 9, + "complexity_size": "M", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 0, + "deterministic_activities": 3, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 3, + "pipeline": "pipeline_web_activity_auth", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 3, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 3, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 6, + "complexity_size": "M", + "control_flow_activities": 3, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 0, + "deterministic_activities": 3, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 0, + "pipeline": "pl_test_appendvariable_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 2, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 2, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 10, + "complexity_size": "M", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 4, + "deterministic_activities": 2, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 2, + "pipeline": "pl_test_copy_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 1, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 1, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 4, + "complexity_size": "S", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 1, + "deterministic_activities": 1, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 1, + "pipeline": "pl_test_delete_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 2, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 2, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 6, + "complexity_size": "M", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 0, + "deterministic_activities": 2, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 2, + "pipeline": "pl_test_executepipeline_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 4, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 4, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 10, + "complexity_size": "M", + "control_flow_activities": 3, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 1, + "deterministic_activities": 4, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 1, + "pipeline": "pl_test_filter_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 2, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 2, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 7, + "complexity_size": "M", + "control_flow_activities": 1, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 2, + "deterministic_activities": 2, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 1, + "pipeline": "pl_test_foreach_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 8, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 8, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 21, + "complexity_size": "L", + "control_flow_activities": 2, + "coverage_pct": 100.0, + "databricks_native_activities": 2, + "datasets": 2, + "deterministic_activities": 8, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 4, + "pipeline": "pl_test_ifcondition_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 2, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 2, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 7, + "complexity_size": "M", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 1, + "deterministic_activities": 2, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 2, + "pipeline": "pl_test_lookup_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 1, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 1, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 2, + "complexity_size": "S", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 1, + "datasets": 0, + "deterministic_activities": 1, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 0, + "pipeline": "pl_test_notebook_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 5, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 5, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 10, + "complexity_size": "M", + "control_flow_activities": 5, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 0, + "deterministic_activities": 5, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 0, + "pipeline": "pl_test_setvariable_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 1, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 1, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 2, + "complexity_size": "S", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 1, + "datasets": 0, + "deterministic_activities": 1, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 0, + "pipeline": "pl_test_sparkjar_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 1, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 1, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 2, + "complexity_size": "S", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 1, + "datasets": 0, + "deterministic_activities": 1, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 0, + "pipeline": "pl_test_sparkpython_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 6, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 6, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 10, + "complexity_size": "M", + "control_flow_activities": 1, + "coverage_pct": 100.0, + "databricks_native_activities": 4, + "datasets": 0, + "deterministic_activities": 6, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 1, + "migration_status": "included", + "other_activities": 1, + "pipeline": "pl_test_switch_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 2, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 2, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 4, + "complexity_size": "S", + "control_flow_activities": 2, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 0, + "deterministic_activities": 2, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 0, + "pipeline": "pl_test_wait_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + }, + { + "activities": 3, + "agentic_activities": 0, + "agentic_provider_version": "", + "agentic_resolution_outcomes": "{}", + "audited_activities": 3, + "code_attached_coverage_pct": 100.0, + "collapsible_patterns": 0, + "complexity_score": 9, + "complexity_size": "M", + "control_flow_activities": 0, + "coverage_pct": 100.0, + "databricks_native_activities": 0, + "datasets": 0, + "deterministic_activities": 3, + "deterministic_coverage_pct": 100.0, + "excluded_activities": 0, + "failed_activities": 0, + "finding_count": 0, + "finding_fingerprints": "[]", + "linked_services": 0, + "migration_status": "included", + "other_activities": 3, + "pipeline": "pl_test_webactivity_coverage", + "reconciliation_status": "not_applicable", + "resolved_agentic_count": 0, + "unresolved_agentic_count": 0, + "unsupported_activities": 0 + } +] \ No newline at end of file diff --git a/tests/unit/test_adf_discovery_mapping.py b/tests/unit/test_adf_discovery_mapping.py new file mode 100644 index 0000000..71d72b3 --- /dev/null +++ b/tests/unit/test_adf_discovery_mapping.py @@ -0,0 +1,423 @@ +"""Tests for the ADF -> shared discovery AST mapper (:mod:`flowx.sources.adf.discovery_mapping`). + +Proves the mapping is 1:1 and lossless: every activity becomes one discovery +node, the ADF type is retained verbatim as ``native_type`` / ``original_type``, +the source dict is preserved on ``raw``, dependency edges keep all of their +outcome conditions, and control-flow nesting maps to labelled container branches. +""" + +from __future__ import annotations + +from flowx.discovery_inventory import STRATEGY_PROPERTY +from flowx.models.adf_ast import AdfDefinitions +from flowx.models.discovery import ( + SOURCE_ADF, + ContainerNode, + GapNode, + SourceGraph, + SourceNode, +) +from flowx.sources.adf.discovery_mapping import ( + adf_definitions_to_source_graphs, + adf_pipeline_to_source_graph, +) +from flowx.sources.adf.loader import _parse_pipeline_json, _parse_trigger_json + + +def _pipeline(activities: list[dict], **props) -> SourceGraph: + data = {"name": props.pop("name", "pl"), "properties": {"activities": activities, **props}} + return adf_pipeline_to_source_graph(_parse_pipeline_json(data)) + + +def test_activity_maps_1to1_retaining_type_and_raw() -> None: + """Each activity becomes one node with its ADF type and raw dict preserved.""" + raw_activity = { + "name": "Run Notebook", + "type": "DatabricksNotebook", + "typeProperties": {"notebookPath": "/nb"}, + } + graph = _pipeline([raw_activity]) + + assert len(graph.tasks) == 1 + node = graph.tasks[0] + assert isinstance(node, SourceNode) + assert node.source == SOURCE_ADF + assert node.name == "Run Notebook" + assert node.task_key == "Run Notebook" + # ADF type is retained verbatim -- native_type is the original_type source. + assert node.native_type == "DatabricksNotebook" + # Verbatim source dict preserved for lossless fallback. + assert node.raw == raw_activity + # Target-side strategy stashed in the properties seam (not a typed field). + assert node.properties[STRATEGY_PROPERTY] == "deterministic" + + +def test_dependencies_capture_all_conditions() -> None: + """A dependency edge keeps every outcome condition, not just the first.""" + activities = [ + {"name": "A", "type": "Wait", "typeProperties": {"waitTimeInSeconds": 1}}, + { + "name": "B", + "type": "DatabricksNotebook", + "dependsOn": [{"activity": "A", "dependencyConditions": ["Succeeded", "Skipped"]}], + }, + ] + graph = _pipeline(activities) + + node_b = next(node for node in graph.tasks if node.name == "B") + assert len(node_b.dependencies) == 1 + assert node_b.dependencies[0].upstream == "A" + assert node_b.dependencies[0].conditions == ["Succeeded", "Skipped"] + + +def test_no_motif_collapse_every_activity_is_a_node() -> None: + """Activities that a motif would merge each stay their own node (no collapse).""" + activities = [ + {"name": "Load", "type": "Copy"}, + {"name": "Notify", "type": "WebActivity", "dependsOn": [{"activity": "Load"}]}, + ] + graph = _pipeline(activities) + assert [node.name for node in graph.tasks] == ["Load", "Notify"] + + +def test_control_flow_maps_to_container_branches() -> None: + """An IfCondition maps to a ContainerNode with true/false branches nested.""" + activities = [ + { + "name": "Check", + "type": "IfCondition", + "typeProperties": { + "ifTrueActivities": [{"name": "T", "type": "DatabricksNotebook"}], + "ifFalseActivities": [{"name": "F", "type": "Wait"}], + }, + } + ] + graph = _pipeline(activities) + + container = graph.tasks[0] + assert isinstance(container, ContainerNode) + assert container.native_type == "IfCondition" + assert list(container.branches.keys()) == ["true", "false"] + assert container.branches["true"][0].name == "T" + assert container.branches["false"][0].name == "F" + + +def test_switch_maps_cases_and_default_to_branches() -> None: + """A Switch maps each case value plus default to its own labelled branch.""" + activities = [ + { + "name": "Route", + "type": "Switch", + "typeProperties": { + "cases": [ + {"value": "gold", "activities": [{"name": "G", "type": "Copy"}]}, + {"value": "silver", "activities": [{"name": "S", "type": "Copy"}]}, + ], + "defaultActivities": [{"name": "D", "type": "Wait"}], + }, + } + ] + graph = _pipeline(activities) + + container = graph.tasks[0] + assert isinstance(container, ContainerNode) + assert list(container.branches.keys()) == ["gold", "silver", "default"] + assert container.branches["gold"][0].name == "G" + assert container.branches["default"][0].name == "D" + + +def test_unsupported_activity_becomes_gap_node() -> None: + """An unsupported ADF type maps to a GapNode carrying the reason and raw.""" + activities = [{"name": "Weird", "type": "TotallyUnknownType"}] + graph = _pipeline(activities) + + node = graph.tasks[0] + assert isinstance(node, GapNode) + assert node.reason is not None and "TotallyUnknownType" in node.reason + assert node.raw == {"name": "Weird", "type": "TotallyUnknownType"} + assert node.properties[STRATEGY_PROPERTY] == "unsupported" + + +def test_policy_maps_retry_and_preserves_timeout_verbatim() -> None: + """Retry count/interval map to typed fields; the ISO timeout rides in extensions.""" + activities = [ + { + "name": "Copy", + "type": "Copy", + "policy": { + "timeout": "0.12:00:00", + "retry": 3, + "retryIntervalInSeconds": 30, + "secureInput": True, + }, + } + ] + graph = _pipeline(activities) + + policy = graph.tasks[0].policy + assert policy is not None + assert policy.max_retries == 3 + assert policy.retry_interval_seconds == 30 + # Timeout normalisation to seconds is a target concern -> preserved verbatim. + assert policy.extensions["timeout"] == "0.12:00:00" + assert policy.extensions["secure_input"] is True + + +def test_graph_carries_parameters_variables_tags_and_folder() -> None: + """Graph-level metadata maps onto the shared fields / properties seam.""" + data = { + "name": "pl", + "properties": { + "activities": [{"name": "N", "type": "DatabricksNotebook"}], + "parameters": {"env": {"type": "String", "defaultValue": "dev"}}, + "variables": {"counter": {"type": "Integer"}}, + "annotations": ["team-a", "prod"], + "folder": {"name": "ingest/bronze"}, + }, + } + graph = adf_pipeline_to_source_graph(_parse_pipeline_json(data)) + + assert graph.source == SOURCE_ADF + assert graph.parameters["env"].type == "String" + assert graph.parameters["env"].default == "dev" + assert graph.variables["counter"].type == "Integer" + assert graph.tags == ["team-a", "prod"] + assert graph.properties["folder"] == "ingest/bronze" + assert graph.raw is not None # verbatim pipeline dict preserved + + +def test_definitions_map_preserves_pipeline_order(adf_definitions) -> None: + """The definitions-level mapper yields one graph per pipeline, in order.""" + graphs = adf_definitions_to_source_graphs(adf_definitions) + assert [g.name for g in graphs] == [p.name for p in adf_definitions.pipelines] + assert all(g.source == SOURCE_ADF for g in graphs) + + +# --------------------------------------------------------------------------- +# Triggers / schedules (BLOCKING 1) +# --------------------------------------------------------------------------- + + +def _definitions_with_trigger(trigger: dict) -> AdfDefinitions: + pipeline = _parse_pipeline_json({"name": "pl_sched", "properties": {"activities": []}}) + return AdfDefinitions(pipelines=[pipeline], triggers=[_parse_trigger_json(trigger)]) + + +def test_schedule_trigger_lands_in_source_graph_schedule() -> None: + """A ScheduleTrigger referencing a pipeline populates that graph's schedule.""" + trigger = { + "name": "tr_daily", + "properties": { + "type": "ScheduleTrigger", + "typeProperties": { + "recurrence": {"frequency": "Day", "interval": 1, "timeZone": "UTC"}, + }, + "pipelines": [{"pipelineReference": {"referenceName": "pl_sched", "type": "PipelineReference"}}], + }, + } + graphs = adf_definitions_to_source_graphs(_definitions_with_trigger(trigger)) + + schedule = graphs[0].schedule + assert schedule is not None + assert schedule.kind == "schedule" + # Recurrence payload preserved verbatim as the expression. + assert schedule.expression == {"frequency": "Day", "interval": 1, "timeZone": "UTC"} + assert schedule.timezone == "UTC" + # Full trigger properties preserved losslessly in extensions. + assert schedule.extensions["trigger_name"] == "tr_daily" + assert schedule.extensions["trigger_type"] == "ScheduleTrigger" + assert "properties" in schedule.extensions + + +def test_tumbling_window_trigger_expression_from_type_properties() -> None: + """A TumblingWindowTrigger keeps its typeProperties as the expression.""" + trigger = { + "name": "tr_tumble", + "properties": { + "type": "TumblingWindowTrigger", + "typeProperties": {"frequency": "Hour", "interval": 1, "startTime": "2024-01-01T00:00:00Z"}, + "pipelines": [{"pipelineReference": {"referenceName": "pl_sched"}}], + }, + } + graphs = adf_definitions_to_source_graphs(_definitions_with_trigger(trigger)) + + schedule = graphs[0].schedule + assert schedule is not None + assert schedule.kind == "interval" + assert schedule.expression == {"frequency": "Hour", "interval": 1, "startTime": "2024-01-01T00:00:00Z"} + + +def test_unreferenced_pipeline_has_no_schedule() -> None: + """A pipeline no trigger references keeps ``schedule is None``.""" + trigger = { + "name": "tr_other", + "properties": { + "type": "ScheduleTrigger", + "typeProperties": {"recurrence": {"frequency": "Day", "interval": 1}}, + "pipelines": [{"pipelineReference": {"referenceName": "some_other_pipeline"}}], + }, + } + graphs = adf_definitions_to_source_graphs(_definitions_with_trigger(trigger)) + assert graphs[0].schedule is None + + +def test_multiple_triggers_preserve_extras_in_extensions() -> None: + """A second trigger for the same pipeline is preserved, not overwritten.""" + pipeline = _parse_pipeline_json({"name": "pl_sched", "properties": {"activities": []}}) + triggers = [ + _parse_trigger_json( + { + "name": "tr_first", + "properties": { + "type": "ScheduleTrigger", + "typeProperties": {"recurrence": {"frequency": "Day", "interval": 1}}, + "pipelines": [{"pipelineReference": {"referenceName": "pl_sched"}}], + }, + } + ), + _parse_trigger_json( + { + "name": "tr_second", + "properties": { + "type": "ScheduleTrigger", + "typeProperties": {"recurrence": {"frequency": "Hour", "interval": 6}}, + "pipelines": [{"pipelineReference": {"referenceName": "pl_sched"}}], + }, + } + ), + ] + graphs = adf_definitions_to_source_graphs(AdfDefinitions(pipelines=[pipeline], triggers=triggers)) + + schedule = graphs[0].schedule + assert schedule is not None + assert schedule.extensions["trigger_name"] == "tr_first" # first wins the typed slot + additional = schedule.extensions["additional_triggers"] + assert len(additional) == 1 + # The additional trigger retains its NAME (not just properties). + assert additional[0]["trigger_name"] == "tr_second" + assert additional[0]["trigger_type"] == "ScheduleTrigger" + assert additional[0]["properties"]["typeProperties"]["recurrence"]["interval"] == 6 + + +def test_triggers_do_not_leak_across_pipelines() -> None: + """A ScheduleSpec is per-pipeline: a later A-only trigger must not appear on B. + + Guards the shared-instance aliasing bug -- trigger_ab references A and B, then + trigger_a references only A. B must keep exactly trigger_ab and gain nothing + from trigger_a. + """ + pipeline_a = _parse_pipeline_json({"name": "pl_a", "properties": {"activities": []}}) + pipeline_b = _parse_pipeline_json({"name": "pl_b", "properties": {"activities": []}}) + triggers = [ + _parse_trigger_json( + { + "name": "tr_ab", + "properties": { + "type": "ScheduleTrigger", + "typeProperties": {"recurrence": {"frequency": "Day", "interval": 1}}, + "pipelines": [ + {"pipelineReference": {"referenceName": "pl_a"}}, + {"pipelineReference": {"referenceName": "pl_b"}}, + ], + }, + } + ), + _parse_trigger_json( + { + "name": "tr_a_only", + "properties": { + "type": "ScheduleTrigger", + "typeProperties": {"recurrence": {"frequency": "Hour", "interval": 2}}, + "pipelines": [{"pipelineReference": {"referenceName": "pl_a"}}], + }, + } + ), + ] + graphs = { + graph.name: graph + for graph in adf_definitions_to_source_graphs( + AdfDefinitions(pipelines=[pipeline_a, pipeline_b], triggers=triggers) + ) + } + + schedule_a = graphs["pl_a"].schedule + schedule_b = graphs["pl_b"].schedule + assert schedule_a is not None and schedule_b is not None + assert schedule_a is not schedule_b # distinct instances, no aliasing + # A picked up the second trigger; B must NOT have leaked it. + assert schedule_a.extensions["additional_triggers"][0]["trigger_name"] == "tr_a_only" + assert "additional_triggers" not in schedule_b.extensions + + +def test_fixture_scheduled_pipeline_gets_schedule(adf_definitions) -> None: + """End-to-end over the fixtures: a trigger-referenced pipeline gets a schedule.""" + graphs = {graph.name: graph for graph in adf_definitions_to_source_graphs(adf_definitions)} + # tr_daily_schedule references pipeline_copy_csv_to_delta in the fixtures. + assert graphs["pipeline_copy_csv_to_delta"].schedule is not None + + +# --------------------------------------------------------------------------- +# Empty / one-sided control flow (BLOCKING 2) +# --------------------------------------------------------------------------- + + +def test_empty_if_condition_stays_a_container_with_both_branches() -> None: + """An IfCondition with no children still maps to a ContainerNode, both branches present.""" + graph = _pipeline([{"name": "Gate", "type": "IfCondition", "typeProperties": {}}]) + + node = graph.tasks[0] + assert isinstance(node, ContainerNode) + assert node.native_type == "IfCondition" + assert list(node.branches.keys()) == ["true", "false"] + assert node.branches["true"] == [] + assert node.branches["false"] == [] + + +def test_empty_for_each_stays_a_container_with_body_branch() -> None: + """A ForEach with no children still maps to a ContainerNode with an empty body.""" + graph = _pipeline([{"name": "Loop", "type": "ForEach", "typeProperties": {}}]) + + node = graph.tasks[0] + assert isinstance(node, ContainerNode) + assert list(node.branches.keys()) == ["body"] + assert node.branches["body"] == [] + + +def test_empty_until_stays_a_container_with_body_branch() -> None: + """An Until with no children still maps to a ContainerNode with an empty body.""" + graph = _pipeline([{"name": "Retry", "type": "Until", "typeProperties": {}}]) + + node = graph.tasks[0] + assert isinstance(node, ContainerNode) + assert node.native_type == "Until" + assert list(node.branches.keys()) == ["body"] + assert node.branches["body"] == [] + + +def test_one_sided_if_keeps_empty_false_branch() -> None: + """An If with only a true branch keeps its false branch present-but-empty.""" + graph = _pipeline( + [ + { + "name": "Gate", + "type": "IfCondition", + "typeProperties": {"ifTrueActivities": [{"name": "T", "type": "Wait"}]}, + } + ] + ) + + node = graph.tasks[0] + assert isinstance(node, ContainerNode) + assert [child.name for child in node.branches["true"]] == ["T"] + assert "false" in node.branches # empty branch is present, not dropped + assert node.branches["false"] == [] + + +def test_empty_switch_stays_a_container_with_default_branch() -> None: + """A Switch with no cases still maps to a ContainerNode with an empty default.""" + graph = _pipeline([{"name": "Route", "type": "Switch", "typeProperties": {}}]) + + node = graph.tasks[0] + assert isinstance(node, ContainerNode) + assert list(node.branches.keys()) == ["default"] + assert node.branches["default"] == [] diff --git a/tests/unit/test_adf_inventory_superset.py b/tests/unit/test_adf_inventory_superset.py new file mode 100644 index 0000000..d8b2f39 --- /dev/null +++ b/tests/unit/test_adf_inventory_superset.py @@ -0,0 +1,125 @@ +"""Consumer-safety tests for the ADF inventory emitted via the shared AST + emitter. + +Two guarantees: + +* **Superset** -- the new ``inventory.json`` keeps every key the historical shape + had (per-activity ``name`` / ``type`` / ``strategy`` / ``depends_on`` and the + ``summary`` count block) byte-compatibly, and only *adds* fields on top. The + historical values are reconstructed here from :func:`build_inventory` (the + authoritative classifier, untouched by this slice). +* **Golden coverage** -- feeding the new inventory through the real consumer + (:func:`flowx.reporting.coverage.build_coverage_rows`) reproduces a committed + snapshot, so downstream reporting is provably unchanged. +""" + +from __future__ import annotations + +import json +from pathlib import Path + +from flowx.reporting.coverage import build_coverage_rows +from flowx.sources.adf.loader import build_inventory, load_adf_definitions, main + +FIXTURES_DIR = Path(__file__).parent.parent / "resources" / "json" +GOLDEN_COVERAGE = Path(__file__).parent.parent / "resources" / "golden" / "adf_fixture_coverage.json" + + +def _run_discover(tmp_path: Path) -> Path: + """Run the ADF discover entry point against the fixtures; return metadata dir.""" + exit_code = main(["--source-dir", str(FIXTURES_DIR), "--output-dir", str(tmp_path)]) + assert exit_code == 0 + return tmp_path / "metadata" + + +def _legacy_pipeline_map() -> dict[str, list[dict]]: + """Reconstruct the pre-change per-pipeline activity shape from build_inventory. + + This is exactly what the removed ``_inventory_to_dict`` used to emit, rebuilt + from the authoritative classifier so the superset check does not depend on a + frozen copy of the old serializer. + """ + inventory = build_inventory(load_adf_definitions(FIXTURES_DIR)) + pipeline_map: dict[str, list[dict]] = {} + for item in inventory.items: + entry: dict = {"name": item.activity_name, "type": item.activity_type, "strategy": item.strategy.value} + if item.depends_on: + entry["depends_on"] = item.depends_on + pipeline_map.setdefault(item.pipeline_name, []).append(entry) + return pipeline_map + + +def test_inventory_top_level_is_superset(tmp_path: Path) -> None: + """Top-level keys include the legacy set plus the new ``source`` discriminator.""" + metadata = _run_discover(tmp_path) + inventory = json.loads((metadata / "inventory.json").read_text()) + + # Legacy top-level keys still present. + for key in ("source_dir", "pipelines", "summary"): + assert key in inventory + # Additive discriminator. + assert inventory["source"] == "adf" + + +def test_summary_counts_match_legacy_classifier(tmp_path: Path) -> None: + """The summary count block is byte-compatible with the legacy classifier.""" + metadata = _run_discover(tmp_path) + summary = json.loads((metadata / "inventory.json").read_text())["summary"] + + legacy = build_inventory(load_adf_definitions(FIXTURES_DIR)) + total = legacy.deterministic_count + legacy.agentic_count + legacy.unsupported_count + assert summary["pipeline_count"] == legacy.pipeline_count + assert summary["activity_count"] == total + assert summary["deterministic_count"] == legacy.deterministic_count + assert summary["agentic_count"] == legacy.agentic_count + assert summary["unsupported_count"] == legacy.unsupported_count + assert summary["coverage_pct"] == round((legacy.deterministic_count + legacy.agentic_count) / total * 100, 1) + + +def test_activity_entries_superset_legacy_shape(tmp_path: Path) -> None: + """Every activity keeps its legacy keys/values and only gains additive fields.""" + metadata = _run_discover(tmp_path) + inventory = json.loads((metadata / "inventory.json").read_text()) + legacy_map = _legacy_pipeline_map() + + # Same pipeline membership (ADF omits zero-activity pipelines; fixtures have none empty). + new_names = [pipeline["name"] for pipeline in inventory["pipelines"]] + assert sorted(new_names) == sorted(legacy_map) + + for pipeline in inventory["pipelines"]: + legacy_entries = legacy_map[pipeline["name"]] + new_entries = pipeline["activities"] + assert len(new_entries) == len(legacy_entries) + for legacy_entry, new_entry in zip(legacy_entries, new_entries): + # Every legacy key/value survives byte-for-byte. + for key, value in legacy_entry.items(): + assert new_entry[key] == value, (pipeline["name"], key) + # Additive standardized fields are present. + assert new_entry["original_type"] == legacy_entry["type"] + assert "dependencies" in new_entry + assert "raw" in new_entry + + +def test_dependencies_field_carries_conditions(tmp_path: Path) -> None: + """The additive ``dependencies`` field carries upstream + conditions per edge.""" + metadata = _run_discover(tmp_path) + inventory = json.loads((metadata / "inventory.json").read_text()) + + # The all-dependency-conditions fixture exercises non-default conditions. + pipeline = next(p for p in inventory["pipelines"] if p["name"] == "pipeline_all_dependency_conditions") + edges = [dependency for activity in pipeline["activities"] for dependency in activity["dependencies"]] + assert edges, "expected at least one dependency edge" + for edge in edges: + assert set(edge.keys()) == {"upstream", "conditions", "resolved"} + assert isinstance(edge["conditions"], list) + + +def test_coverage_output_matches_golden(tmp_path: Path) -> None: + """The real coverage consumer reproduces the committed golden snapshot. + + This is the no-consumer-breakage guarantee: regenerate the golden with + ``make test`` after an intentional change and review the diff. + """ + metadata = _run_discover(tmp_path) + rows = json.loads(json.dumps(build_coverage_rows(metadata), sort_keys=True)) + golden = json.loads(GOLDEN_COVERAGE.read_text()) + assert rows == golden diff --git a/tests/unit/test_discovery_inventory.py b/tests/unit/test_discovery_inventory.py new file mode 100644 index 0000000..bb95768 --- /dev/null +++ b/tests/unit/test_discovery_inventory.py @@ -0,0 +1,168 @@ +"""Tests for the source-agnostic inventory emitter (:mod:`flowx.discovery_inventory`). + +These tests build the shared discovery AST by hand -- no ADF, no Airflow -- so +they prove the emitter is genuinely source-agnostic: it takes ``SourceGraph`` +objects in and projects the ``inventory.json`` document out, with no coupling to +any particular front-end. +""" + +from __future__ import annotations + +import ast + +import flowx.discovery_inventory as discovery_inventory +from flowx.discovery_inventory import STRATEGY_PROPERTY, build_source_inventory +from flowx.models.discovery import ( + CONCEPT_BRANCH, + CONCEPT_NOTEBOOK, + ContainerNode, + SourceDependency, + SourceGraph, + SourceNode, +) + + +def _node(task_key: str, native_type: str, strategy: str, *, deps: list[SourceDependency] | None = None) -> SourceNode: + return SourceNode( + source_id=task_key, + task_key=task_key, + concept=CONCEPT_NOTEBOOK, + source="unit", + name=task_key, + native_type=native_type, + dependencies=deps or [], + properties={STRATEGY_PROPERTY: strategy}, + raw={"name": task_key, "type": native_type}, + ) + + +def test_emitter_has_no_source_specific_imports() -> None: + """The emitter module must not import any per-source package. + + Source-agnostic means the ADF/Airflow loaders depend on the emitter, never + the other way round. Guard that by inspecting the module's actual import + statements (not arbitrary text -- the docstring legitimately names the + ``"adf"`` / ``"airflow"`` discriminator values). + """ + source = (discovery_inventory.__file__ or "").rstrip("c") + with open(source, encoding="utf-8") as handle: + tree = ast.parse(handle.read()) + + imported: list[str] = [] + for node in ast.walk(tree): + if isinstance(node, ast.Import): + imported.extend(alias.name for alias in node.names) + elif isinstance(node, ast.ImportFrom) and node.module: + imported.append(node.module) + + assert not any(name.startswith("flowx.sources") for name in imported), imported + + +def test_top_level_shape_and_summary_counts() -> None: + """A hand-built graph projects to the canonical top-level shape and counts.""" + graph = SourceGraph( + name="g1", + source="unit", + tasks=[ + _node("n1", "Notebook", "deterministic"), + _node("n2", "DataFlow", "agentic"), + _node("n3", "Mystery", "unsupported"), + ], + ) + + inventory = build_source_inventory([graph], source="unit", source_dir="/tmp/src") + + assert sorted(inventory.keys()) == ["pipelines", "source", "source_dir", "summary"] + assert inventory["source"] == "unit" + assert inventory["source_dir"] == "/tmp/src" + assert inventory["summary"] == { + "pipeline_count": 1, + "activity_count": 3, + "deterministic_count": 1, + "agentic_count": 1, + "unsupported_count": 1, + "coverage_pct": 66.7, + } + + +def test_activity_entry_carries_legacy_and_additive_fields() -> None: + """Each activity keeps the legacy keys and gains the additive standardized ones.""" + graph = SourceGraph( + name="g1", + source="unit", + tasks=[ + _node("a", "Notebook", "deterministic"), + _node( + "b", + "Copy", + "deterministic", + deps=[SourceDependency(upstream="a", conditions=["Succeeded", "Skipped"])], + ), + ], + ) + + inventory = build_source_inventory([graph], source="unit", source_dir="/tmp/src") + entries = {entry["name"]: entry for entry in inventory["pipelines"][0]["activities"]} + + # Legacy keys (byte-compatible with the historical shape). + assert entries["a"]["type"] == "Notebook" + assert entries["a"]["strategy"] == "deterministic" + assert "depends_on" not in entries["a"] # no deps -> key omitted, as before + assert entries["b"]["depends_on"] == ["a"] + + # Additive standardized fields. + assert entries["a"]["original_type"] == "Notebook" + assert entries["a"]["dependencies"] == [] + assert entries["a"]["raw"] == {"name": "a", "type": "Notebook"} + assert entries["b"]["dependencies"] == [{"upstream": "a", "conditions": ["Succeeded", "Skipped"], "resolved": True}] + + +def test_container_branches_are_flattened_in_source_order() -> None: + """Container children are flattened depth-first in branch declaration order.""" + branch_true = _node("t", "Notebook", "deterministic") + branch_false = _node("f", "Notebook", "deterministic") + container = ContainerNode( + source_id="if", + task_key="if", + concept=CONCEPT_BRANCH, + source="unit", + name="if", + native_type="IfCondition", + properties={STRATEGY_PROPERTY: "deterministic"}, + branches={"true": [branch_true], "false": [branch_false]}, + ) + graph = SourceGraph(name="g", source="unit", tasks=[container, _node("after", "Notebook", "deterministic")]) + + inventory = build_source_inventory([graph], source="unit", source_dir="/tmp/src") + names = [entry["name"] for entry in inventory["pipelines"][0]["activities"]] + + assert names == ["if", "t", "f", "after"] + assert inventory["summary"]["activity_count"] == 4 + + +def test_include_empty_pipelines_toggle_preserves_membership_semantics() -> None: + """Zero-activity graphs stay counted in summary but drop from the listing when asked. + + Reproduces ADF's long-standing rule: a pipeline with no activities is omitted + from ``pipelines`` yet still counted in ``summary.pipeline_count``. + """ + empty = SourceGraph(name="empty", source="unit", tasks=[]) + populated = SourceGraph(name="full", source="unit", tasks=[_node("n", "Notebook", "deterministic")]) + + omitted = build_source_inventory( + [empty, populated], source="unit", source_dir="/tmp", include_empty_pipelines=False + ) + assert [p["name"] for p in omitted["pipelines"]] == ["full"] + assert omitted["summary"]["pipeline_count"] == 2 # empty still counted + + listed = build_source_inventory([empty, populated], source="unit", source_dir="/tmp", include_empty_pipelines=True) + assert [p["name"] for p in listed["pipelines"]] == ["empty", "full"] + assert listed["summary"]["pipeline_count"] == 2 + + +def test_empty_input_yields_zero_coverage() -> None: + """No graphs -> empty listing, zeroed summary, 0.0 coverage (no divide-by-zero).""" + inventory = build_source_inventory([], source="unit", source_dir="/tmp") + assert inventory["pipelines"] == [] + assert inventory["summary"]["coverage_pct"] == 0.0 + assert inventory["summary"]["pipeline_count"] == 0