Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
eb768bb
ADF discovery: inventory + lineage + motifs on the shared source AST
matthewmoorcroft Sep 16, 2026
a8aa8e9
Discover: loud 0-pipeline warning + self-describing activity-count units
matthewmoorcroft Sep 18, 2026
79c41dd
Discover 0-pipeline warning: name the user-facing --adf-source-path flag
matthewmoorcroft Sep 18, 2026
e84c222
Merge discovery/shared-ast-standard (shared emitter, source-graph env…
matthewmoorcroft Oct 5, 2026
e39a05f
Persist ADF source graphs at discover through the shared envelope
matthewmoorcroft Oct 5, 2026
db46419
Call the ADF mapper's target the discovery graph contract
matthewmoorcroft Oct 5, 2026
22e406c
Merge discovery/shared-ast-standard (motifs in the hashed graph, inve…
matthewmoorcroft Oct 5, 2026
b9e3305
Save ADF motifs in the source graphs and project the inventory from t…
matthewmoorcroft Oct 5, 2026
770c0b2
Merge discovery/shared-ast-standard (validated head: inventory lineag…
matthewmoorcroft Oct 5, 2026
b21c15d
no-mistakes(review): Qualify ADF dataset identities and share the lin…
matthewmoorcroft Oct 6, 2026
1068dce
no-mistakes(review): Reject parameterised linked-service parts in ADF…
matthewmoorcroft Oct 6, 2026
e50cf0f
no-mistakes(document): Document ADF source_graphs.json and inventory …
matthewmoorcroft Oct 6, 2026
1ded91c
Tighten ADF lineage identities and skip external run-now edges
matthewmoorcroft Oct 7, 2026
bdf0cee
Merge discovery/shared-ast-standard (#84 validated head bbd00f1) into…
matthewmoorcroft Oct 7, 2026
d64624c
no-mistakes(review): Keep bracket-quoted SQL names physical; reword g…
matthewmoorcroft Oct 7, 2026
adf84b5
no-mistakes(review): Drop identities for storeSettings-overridden rea…
matthewmoorcroft Oct 7, 2026
a911dc7
no-mistakes(review): Drop identities for query, procedure and path-ov…
matthewmoorcroft Oct 7, 2026
a4bae0e
no-mistakes(review): Match source query and procedure overrides by ke…
matthewmoorcroft Oct 7, 2026
79f70e3
no-mistakes(document): Document run-now skip in build_control_edges d…
matthewmoorcroft Oct 7, 2026
efd8208
Drop signatures on overridden reads, keep empty Switch cases, carry m…
matthewmoorcroft Oct 7, 2026
bb6763e
no-mistakes(review): Key Switch cases like parser, fix override docst…
matthewmoorcroft Oct 7, 2026
7919ce2
no-mistakes(document): Docs already current; lint and type checks clean
matthewmoorcroft Oct 7, 2026
0c6c1bd
Remove a stray scratch file the document step committed
matthewmoorcroft Oct 7, 2026
820b48b
no-mistakes(review): Key Switch case valued "default" as case:default
matthewmoorcroft Oct 7, 2026
edc0756
no-mistakes(review): Document case:default key in ContainerNode contract
matthewmoorcroft Oct 7, 2026
cac27d8
no-mistakes(test): Keep Switch "default" case label unique against de…
matthewmoorcroft Oct 7, 2026
3e51c68
Qualify the weak path signature with the dataset's literal store
matthewmoorcroft Oct 8, 2026
a5b3526
Keep the ADF case:default rule out of the source-neutral ContainerNod…
matthewmoorcroft Oct 8, 2026
13d16b7
no-mistakes(review): Qualify hidden-account file signatures with link…
matthewmoorcroft Oct 8, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -56,7 +56,7 @@ All three phases write into one shared `<output_dir>` (default `./flowx_output`)
the DAB bundle at the top level, kept artifacts under `metadata/`, and transient
intermediates under `.work/` (pruned by `package`).

1. **Discover** -- Parse ADF JSON from UC volumes -> typed AST -> `metadata/inventory.json` + `metadata/profile_report.csv` + verbatim `metadata/<pipeline>.arm.json`
1. **Discover** -- Parse ADF JSON from UC volumes -> typed AST -> `metadata/source_graphs.json` (hashed) -> `metadata/inventory.json` + `metadata/profile_report.csv` + verbatim `metadata/<pipeline>.arm.json`
2. **Convert** -- Registry dispatch + topological sort -> Pipeline IR (deterministic + agentic gaps); transient report at `.work/translation_report.json`
3. **Package** -- IR -> DAB YAML + generated notebooks + setup scripts; prunes `.work/`

Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -257,6 +257,7 @@ flowx_output/
SETUP.md # Setup instructions (package)
metadata/
inventory.json # discover: activity inventory
source_graphs.json # discover (ADF): saved source graphs the inventory is built from
profile_report.csv # discover: per-pipeline complexity report
<pipeline>.arm.json # discover: verbatim original ADF/ARM source
configuration.json # modify: collected configuration answers
Expand Down
1 change: 1 addition & 0 deletions skills/flowx-discover/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,7 @@ All under the shared `<output_dir>/metadata/` folder:
| File | Description |
|---|---|
| `metadata/inventory.json` | Classified activity inventory for the convert phase |
| `metadata/source_graphs.json` | (ADF) Saved, content-hashed source graphs (activities, lineage, detected motifs); the inventory is built from it and records its hash as `source_graphs_sha256` |
| `metadata/profile_report.csv` | Per-pipeline complexity report (counts + T-shirt size) |
| `metadata/<pipeline>.arm.json` | (ADF) Verbatim original source for each pipeline (provenance) |

Expand Down
14 changes: 11 additions & 3 deletions skills/flowx-discover/sources/adf.md
Original file line number Diff line number Diff line change
Expand Up @@ -51,16 +51,24 @@ Read `<output_dir>/metadata/inventory.json`:
{
"name": "PipelineName",
"activities": [
{"name": "CopyFromBlob", "type": "Copy", "strategy": "deterministic", "translator": "copy.py"},
{"name": "RunDataFlow", "type": "ExecuteDataFlow", "strategy": "agentic"}
{"name": "CopyFromBlob", "type": "Copy", "strategy": "deterministic", "task_key": "CopyFromBlob"},
{"name": "RunDataFlow", "type": "ExecuteDataFlow", "strategy": "agentic", "task_key": "RunDataFlow"}
]
}
],
"summary": {"pipeline_count": 12, "activity_count": 47, "deterministic_count": 35,
"agentic_count": 10, "unsupported_count": 2, "coverage_pct": 95.7}
"agentic_count": 10, "unsupported_count": 2, "coverage_pct": 95.7},
"source_graphs_sha256": "<document_sha256 of metadata/source_graphs.json>"
}
```

Each activity's `task_key` equals its ADF activity name. `source_graphs_sha256` names the saved
`metadata/source_graphs.json` the inventory was built from.

The inventory deliberately has no top-level `generated_at` timestamp any more, so the same export
always produces byte-identical `inventory.json` and any hash taken over it stays stable. Use the
file's modification time if you need to know when discover ran.

## Step 4b — Review the complexity report

`<output_dir>/metadata/profile_report.csv` has one row per pipeline: `pipeline`, `activities`,
Expand Down
9 changes: 8 additions & 1 deletion src/flowx/bundler/dab_writer.py
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@
from flowx.bundler.notebook_writer import write_notebooks
from flowx.bundler.prereqs_writer import ManualParameter, build_prereqs, render_setup_md
from flowx.bundler.setup_generator import generate_setup_tasks
from flowx.ir_serde import data_asset_from_dict, lineage_from_dict
from flowx.models.dab import DabNotebook
from flowx.models.ir import (
Activity,
Expand Down Expand Up @@ -2021,6 +2022,7 @@ def pipeline_dict_to_ir(pipeline_dict: dict[str, Any]) -> tuple[Pipeline, list[d
):
raise ValueError(f"Invalid Pipeline email notification entry: {event!r}")
email_notifications[str(event)] = list(recipients)
raw_lineage = pipeline_dict.get("lineage")

pipeline = Pipeline(
name=pipeline_dict.get("name", "unknown"),
Expand All @@ -2037,6 +2039,7 @@ def pipeline_dict_to_ir(pipeline_dict: dict[str, Any]) -> tuple[Pipeline, list[d
reconciliation_status=pipeline_dict.get("reconciliation_status"),
migration_status=pipeline_dict.get("migration_status", "included"),
audit=dict(pipeline_dict.get("audit") or {}),
lineage=lineage_from_dict(raw_lineage) if raw_lineage else None,
)
return pipeline, parameters

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


Expand Down
52 changes: 4 additions & 48 deletions src/flowx/discovery_serde.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,8 @@
resolvable physical ``identity`` when there is one (else ``None``), an
always-present ``signature``, and an open ``asset_type`` that also covers
non-physical / logical / value hand-offs (e.g. an Airflow XCom). 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.
:class:`~flowx.models.ir.Lineage` block goes through ``ir_serde``'s
``lineage_to_dict`` / ``lineage_from_dict`` pair for the same reason.

Every node dict carries a ``node_type`` discriminator (the dataclass name) so a
:class:`~flowx.models.discovery.ContainerNode` or
Expand All @@ -33,7 +32,7 @@
from pathlib import Path
from typing import Any

from flowx.ir_serde import data_asset_from_dict, data_asset_to_dict, lineage_to_dict
from flowx.ir_serde import data_asset_from_dict, data_asset_to_dict, lineage_from_dict, lineage_to_dict
from flowx.models.discovery import (
ContainerNode,
GapNode,
Expand All @@ -44,7 +43,6 @@
SourceGraph,
SourceNode,
)
from flowx.models.ir import ControlEdge, DataEdge, Lineage, MotifAnnotation

SOURCE_GRAPHS_FILENAME = "source_graphs.json"
SOURCE_GRAPHS_CONTRACT_VERSION = "1"
Expand Down Expand Up @@ -97,7 +95,7 @@ def source_graph_from_dict(raw: dict[str, Any]) -> SourceGraph:
run_timeout_seconds=raw.get("run_timeout_seconds"),
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,
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"),
Expand Down Expand Up @@ -188,48 +186,6 @@ def _canonical_sha256(value: Any) -> str:
return hashlib.sha256(canonical.encode("utf-8")).hexdigest()


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 []),
source_type_hint=motif.get("source_type_hint"),
)
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:
Expand Down
54 changes: 53 additions & 1 deletion src/flowx/ir_serde.py
Original file line number Diff line number Diff line change
Expand Up @@ -22,8 +22,10 @@
from flowx.models.ir import (
Activity,
AppendVariableActivity,
ControlEdge,
CopyActivity,
DataAsset,
DataEdge,
DbtFactoryActivity,
DeleteActivity,
ExecutePipelineActivity,
Expand Down Expand Up @@ -83,6 +85,8 @@ def pipeline_to_dict(pipeline: Pipeline) -> dict[str, Any]:
}
if pipeline.translation_configuration is not None:
result["translation_configuration"] = configuration_to_dict(pipeline.translation_configuration)
if pipeline.lineage is not None:
result["lineage"] = lineage_to_dict(pipeline.lineage)
return result


Expand Down Expand Up @@ -148,6 +152,48 @@ def lineage_to_dict(lineage: Lineage) -> dict[str, Any]:
}


def lineage_from_dict(raw: dict[str, Any]) -> Lineage:
"""Rehydrate a :class:`Lineage` block from the dict :func:`lineage_to_dict` emits.

The one inverse shared by the translation report and the discovery graph, so a
lineage block reads back the same way wherever it was written.
"""
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 []),
source_type_hint=motif.get("source_type_hint"),
)
for motif in raw.get("motifs") or []
],
)


def _motif_annotation_to_dict(motif: MotifAnnotation) -> dict[str, Any]:
"""Serialise one motif annotation; ``source_type_hint`` is written only when set.

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

extra = activity_extra_fields(task)
task_dict.update(extra)
Expand Down Expand Up @@ -420,7 +472,7 @@ def activity_extra_fields(activity: Activity) -> dict[str, Any]:
if activity.job_parameters:
extra["job_parameters"] = activity.job_parameters
case MotifActivity():
extra["motif_id"] = activity.motif_id
# motif_id now lives on the Activity base and is serialised by activity_to_dict.
extra["display_name"] = activity.display_name
extra["databricks_replacement"] = activity.databricks_replacement
extra["matched_activity_names"] = activity.matched_activity_names
Expand Down
Loading
Loading