Skip to content

ADF: map discovery AST -> IR (carry lineage through conversion) #85

Description

@matthewmoorcroft

Summary: Phase 2 makes ADF conversion decision-fed: consume the persisted shared discovery AST plus the user-approved route plan and build IR directly, including lineage, cross-pipeline motifs, and typed Lakeflow Connect ingestion IR.

Depends on: #84, the main re-land of #79-#81, routing #77, and the Phase 1 validation seams from #87.

Blocks: Completion of the ADF decision-fed conversion architecture and the shared IR surface reused by #86.

Current status

This is an open issue with no Phase 2 implementation yet.

On main today:

  • src/flowx/sources/adf/translate.py:main calls sources/adf/loader.py:load_adf_definitions again during convert, then translate_pipeline dispatches the ADF-specific AST through per-activity translators.
  • translate_pipeline detects motifs after activity translation and may call motifs/collapser.py:collapse_motifs. It does not consume inventory.json, SourceGraph, lineage, insights, or a connected-component decision.
  • ir_serde.py:pipeline_to_dict/activity_to_dict writes the convert-to-package report; bundler/dab_writer.py:pipeline_dict_to_ir/_reconstruct_ir rehydrates it for packaging.
  • Copy-driven Lakeflow Connect is represented as flags on models/ir.py:CopyActivity and materialized by preparer/activity_preparers/copy.py:_prepare_lakeflow_connect_copy. A second path exists only for consolidated metadata-driven MotifActivity in preparer/activity_preparers/motif.py:_prepare_consolidated_metadata_driven. There is no typed ingestion IR independent of Copy/motif collapse.

On integration/discovery, #79-#81 build SourceGraph and emit its projected inventory, but discovery currently persists the inventory projection rather than a complete SourceGraph artifact suitable for direct conversion.

What needs to change

After Phase 1 is stable, make ADF convert consume:

  1. the persisted, source-faithful SourceGraph(s) produced by discovery;
  2. the enhanced inventory (lineage, motifs, insights);
  3. the user-approved connected-component conversion plan from Phase 1: Routing agent — per-connected-component deterministic/agentic decision #77.

Build Pipeline/Activity IR directly from those inputs. Both phases are decision-driven; the Phase 2 difference is that the decision is applied inside conversion, not as the post-convert alteration used by #87.

The mapper must carry node data_reads/data_writes and SourceGraph.lineage, reconcile discovery task keys with normalized IR task keys, preserve lineage through motif transformations, and retain raw/extension information needed by target translators.

Phase 2 also expands deterministic coverage:

  • detect/route cross-pipeline motifs using component context rather than one AdfPipeline at a time;
  • add typed Lakeflow Connect ingestion IR so managed ingestion is not reachable only through CopyActivity flags or a consolidated MotifActivity.

How to approach

  • Persist complete SourceGraph documents during sources/adf/loader.py:main using discovery_serde.py:source_graph_to_dict, under metadata, while keeping inventory.json as the additive consumer-facing projection.
  • Add the inverse loader at the convert boundary using discovery_serde.py:source_graph_from_dict; do not reparse ADF JSON in the normal Phase 2 path.
  • Add an ADF SourceNode/ContainerNode-to-Activity mapper beside sources/adf/discovery_mapping.py, reusing existing translators only at well-defined native_type/raw target-lowering seams.
  • Extend models/ir.py minimally for carried lineage/assets and a typed ingestion activity/resource model. Update ir_serde.py and bundler/dab_writer.py:_reconstruct_ir before fallback.
  • Apply the Phase 1: Routing agent — per-connected-component deterministic/agentic decision #77 decision per connected component during mapping: deterministic components use the typed mapper; agentic decisions use the same in-engine component/report contract established by Phase 1: In-engine agentic escape hatch (passthrough + post-convert IR alteration) #87.
  • Update motifs/detector.py/collapser.py or add a component-level detector so lineage/task keys are remapped rather than dropped when motifs span pipelines.
  • Route typed ingestion through preparer/workflow_preparer.py and existing dab_writer.py pipeline_resources; validate every pipeline_task reference in validate/bundle_invariants.py.

Acceptance/verification:

  • Golden tests prove discovery AST -> IR -> serde -> rehydrate -> package preserves dependencies, policies, schedules, parameters, assets, lineage, and raw provenance.
  • Task-key reconciliation tests prove every lineage endpoint resolves after normalization/collapse.
  • Connected-component fixtures cover deterministic, agentic, and mixed recommendations rejected by plan validation.
  • Typed Lakeflow Connect tests structurally validate generated ingestion resources and pipeline_task references.
  • Existing inventory coverage golden remains green.
  • Default/compatibility tests show Phase 2 can be gated until enabled; Airflow's existing path is untouched.

This is a later, decision-fed conversion change; it must not be folded into the Phase 1 re-land.

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions