diff --git a/.claude-plugin/plugin.json b/.claude-plugin/plugin.json index e41c463..fa3e904 100644 --- a/.claude-plugin/plugin.json +++ b/.claude-plugin/plugin.json @@ -11,7 +11,10 @@ "skills": [ "./skills/flowx-setup", "./skills/flowx-discover", + "./skills/flowx-enrich", + "./skills/flowx-route", "./skills/flowx-convert", + "./skills/flowx-resolve-airflow-gaps", "./skills/flowx-package", "./skills/flowx-migrate" ] diff --git a/skills/flowx-discover/SKILL.md b/skills/flowx-discover/SKILL.md index 4d776bb..047e89d 100644 --- a/skills/flowx-discover/SKILL.md +++ b/skills/flowx-discover/SKILL.md @@ -74,17 +74,28 @@ The inventory classifies every task into one of three strategies: - **Agentic** — requires LLM-assisted translation from the source definition. - **Unsupported** — no known translation path; needs manual intervention. -## Step 3 — Optional: enrich the inventory with agentic insights +## Step 3 — Enrich the inventory (default next step) -After the deterministic inventory is written, you can *author* a layer of judgment the parser -cannot derive — factory-wide architecture recommendation, per-pipeline intent + recommended -Databricks patterns, and cross-pipeline relationships — and merge it back under a single additive -`insights` key. flowx contains no LLM: you author the JSON, the library validates and merges it. -This is optional and changes nothing about conversion. See **`insights.md`** for the shape, the -validation rules, and the `enrich` command (CLI and MCP). +Once the deterministic inventory is written, the standard flow **chains into enrichment**: you author +a layer of judgment the parser cannot derive — a factory-wide architecture recommendation, +per-pipeline intent + recommended Databricks patterns, and cross-pipeline relationships — and merge +it back under a single additive `insights` key. flowx contains no LLM: you author the JSON, the +library validates and merges it. The routing step (`flowx-route`) reads this block to present the +agentic conversion option per pipeline group. + +**Continue with the `flowx-enrich` skill by default** — it guides authoring the insights and running +`enrich`. Enrichment is additive and leaves every deterministic inventory key byte-identical, so it +never destabilizes discover's output. + +**Deterministic-only skip path.** A standalone, headless discover — no agent, no LLM — is fully +valid: `metadata/inventory.json` is complete and self-standing without an `insights` block, and +`flowx-route` still recommends and records a plan from the deterministic structure alone. Skip +`flowx-enrich` only when the caller explicitly asked for a deterministic-only pass, or when no human +reader needs a migration narrative. ## Reference - `sources/adf.md` — Azure Data Factory discovery (ARM JSON, UC-volume download, complexity report) - `sources/airflow.md` — Apache Airflow discovery (DAG `.py` parsing, operator classification) -- `insights.md` — authoring the optional agentic `insights` layer and running `enrich` +- `flowx-enrich` skill — authoring the agentic `insights` layer and running `enrich` (the default + next step) diff --git a/skills/flowx-enrich/SKILL.md b/skills/flowx-enrich/SKILL.md new file mode 100644 index 0000000..e441cba --- /dev/null +++ b/skills/flowx-enrich/SKILL.md @@ -0,0 +1,109 @@ +--- +name: flowx-enrich +description: > + Enrich the discover inventory with an agent-authored layer of judgment — a factory-wide + architecture recommendation, per-pipeline intent + recommended Databricks patterns, and + cross-pipeline relationships — then validate and merge it into inventory.json. The default next + step after flowx-discover and the input the routing step consumes. +triggers: + - "enrich inventory" + - "enrich pipelines" + - "author insights" + - "agentic insights" + - "recommend databricks patterns" + - "annotate inventory" + - "enrich discover" +--- + +# Enrich the Inventory with Agentic Insights + +The deterministic discover pass records what each source workflow **is**; it cannot record what to +**do** about it. That judgment — a factory-wide architectural recommendation, each pipeline's intent +and recommended Databricks patterns, and how pipelines couple — is authored by **you, the agent**, +and merged back into `metadata/inventory.json` under a single additive `insights` key. + +This is the standard step **between discover and route** in the flowx workflow. `flowx-discover` +chains into this skill by default; the routing step (`flowx-route`) reads the `insights` block to +present the agentic conversion option per connected component. + +## No LLM inside flowx — you author, the library validates and merges + +**There is no LLM inside flowx.** You author the insights JSON; the library (`enrich`) only +*validates and merges* it — the same author → validate → merge, fingerprint-bound contract the +agentic gap-resolution and routing paths use. That keeps the deterministic inventory trustworthy and +every insight accountable: foreign keys must point at real pipelines, and every cross-pipeline edge +is either an annotation of a proven lineage edge or an explicitly-flagged inference with cited +evidence. + +`enrich` is **additive**: it merges only an `insights` block and leaves every existing inventory key +byte-identical. It changes no conversion, IR, or routing decision on its own — the `insights` are +descriptive data that `flowx-route` later consumes. + +## When to skip enrich (deterministic-only) + +Enrichment is on by default, but it is skippable. A standalone, deterministic-only discover — no +agent, no LLM — is fully valid: `metadata/inventory.json` from discover is complete and self-standing +without an `insights` block, and `flowx-route` still recommends and records a plan from the +deterministic structure alone (the agentic option simply shows no recommended patterns). Skip enrich +when the caller asked for a headless/deterministic pass, or when no human reader needs a migration +narrative. + +## How to author (three steps) + +1. **Read the deterministic inventory.** Load `/metadata/inventory.json`. Note every + pipeline `name` (these are the only valid foreign keys), and each pipeline's `lineage` block — in + particular `lineage.control_edges`, each `{source_workflow, target_workflow, via_task_key}`. A + deterministic **control** relationship you annotate must match one of these exactly. +2. **Read the source artifacts** you need to form judgment — the per-pipeline `raw` payloads in the + inventory, the ADF `metadata/.arm.json` provenance, or the DAG source — enough to state + each pipeline's *intent* and the Databricks patterns that fit. Ground every recommended pattern in + a **real, publicly-documented** Databricks capability; never invent a product name. +3. **Author the insights JSON, then call `enrich`.** The library validates it against the inventory + and, only when clean, merges it in atomically. On any violation the inventory is left untouched + and you get the full list of problems to fix in one pass. + +See **`insights.md`** in this skill directory for the exact insights shape, every field, and the +validation rules the library enforces. + +## Run enrich — MCP tool or venv CLI + +Run the **`setup`** skill first if you haven't. Both paths run the same validate-and-merge contract. + +- **MCP tool (Databricks Genie Code, or a local stdio registration):** call the single **`flowx`** + tool with `command="enrich"` and either inline insights or a file: + + ``` + flowx(command="enrich", parameters={"output_dir": "", "insights": { ... }}) # inline object + flowx(command="enrich", parameters={"output_dir": "", "insights_path": ""}) + ``` + + Provide **exactly one** of `insights` (inline object) or `insights_path`. `ok` reflects validation; + `result.violations` lists any problems. Run **no** `python3`/`$PY` commands on this path. + +- **venv CLI (local, no MCP server):** ensure the venv exists (`setup` / `bootstrap.sh`), then: + + ```bash + export PYTHONPATH="/src" + PY="$(cat /.migration-venv)" + "$PY" -m flowx.adapter enrich --output-dir --insights-path insights.json + ``` + + Both `--output-dir` and `--insights-path` are required; `--out ` optionally writes the result + JSON to a file instead of stdout. **Exit code 0** means the insights merged; **exit code 1** prints + the violations JSON and leaves the inventory untouched. + +## Idempotency & safety + +`enrich` is atomic and idempotent: it replaces the whole `insights` block (never stacks), recomputes +the `inventory_sha256` fingerprint from the deterministic inventory, and leaves every existing +inventory key byte-identical. Re-running with the same insights rewrites the same bytes; re-running +with different insights replaces the block. A validation failure writes nothing. + +## Next step + +After the inventory is enriched, continue with **`flowx-route`** to recommend and record a +per-connected-component conversion route (deterministic vs. agentic), then convert and package. + +## Reference + +- `insights.md` — the insights shape, every field, and the validator's rules. diff --git a/skills/flowx-discover/insights.md b/skills/flowx-enrich/insights.md similarity index 52% rename from skills/flowx-discover/insights.md rename to skills/flowx-enrich/insights.md index 75508b4..93db744 100644 --- a/skills/flowx-discover/insights.md +++ b/skills/flowx-enrich/insights.md @@ -1,53 +1,13 @@ -# Authoring agentic insights (optional, after discover) +# The agentic insights shape -The deterministic discover pass records what each source workflow **is**; it cannot record what -to **do** about it. That judgment — a factory-wide architectural recommendation, each pipeline's -intent and recommended Databricks patterns, and how pipelines couple — is authored by **you, the -agent**, and merged back into `inventory.json` under a single additive `insights` key. +This is the reference for the `insights` JSON you author before calling `enrich`. The `flowx-enrich` +SKILL.md covers the workflow (no-LLM contract, the three authoring steps, how to run `enrich`, and +when a deterministic-only pass skips it); this file covers **what to write** and the rules the +validator enforces. -**There is no LLM inside flowx.** You author the insights JSON; the library only *validates and -merges* it (the same author-then-validate-merge contract the agentic gap-resolution path uses). -That keeps the deterministic inventory trustworthy and every insight accountable — foreign keys -must point at real pipelines, and every cross-pipeline edge is either an annotation of a proven -lineage edge or an explicitly-flagged inference with cited evidence. - -Insights are **optional** and change nothing about conversion. They are descriptive data that a -later routing step may consume; on their own they add zero routing/IR/conversion decisions. - -## When to author them - -After `discover` has written `metadata/inventory.json`, and before (or instead of) convert, when a -human reader would benefit from a migration narrative: which pipelines collapse onto a managed -capability, how the factory hangs together, what the risky couplings are. - -## How to author (three steps) - -1. **Read the deterministic inventory.** Load `/metadata/inventory.json`. Note every - pipeline `name` (these are the only valid foreign keys), and each pipeline's `lineage` block — - in particular `lineage.control_edges`, each `{source_workflow, target_workflow, via_task_key}`. - A deterministic **control** relationship you annotate must match one of these exactly. -2. **Read the source artifacts** you need to form judgment — the per-pipeline `raw` payloads in the - inventory, the ADF `metadata/.arm.json` provenance, or the DAG source — enough to state - each pipeline's *intent* and the Databricks patterns that fit. Ground every recommended pattern - in a **real, publicly-documented** Databricks capability; never invent a product name. -3. **Author the insights JSON, then call `enrich`.** The library validates it against the inventory - and, only when clean, merges it in atomically. On any violation the inventory is left untouched - and you get the full list of problems to fix in one pass. - -### Run enrich - -- **MCP tool:** `flowx("enrich", {"output_dir": "", "insights": { ... }})` (inline object), or - pass `"insights_path": ""` instead. `ok` reflects validation; `result.violations` lists any - problems. -- **venv CLI:** - - ```bash - export PYTHONPATH="/src" - PY="$(cat /.migration-venv)" - "$PY" -m flowx.adapter enrich --output-dir --insights-path insights.json - ``` - - Exit code 0 means merged; exit code 1 prints the violations JSON and leaves the inventory untouched. +Author the insights when a human reader would benefit from a migration narrative — which pipelines +collapse onto a managed capability, how the factory hangs together, what the risky couplings are. +The `flowx-route` step reads this block to present the agentic conversion option per component. ## The insights shape @@ -103,11 +63,13 @@ fingerprint that binds your insights to the exact inventory they describe): - **Foreign keys.** Every `pipeline_insights[].pipeline` and every relationship `from_pipeline` / `to_pipeline` must be a real pipeline name in the inventory. -- **`recommended_patterns`** (per pipeline and system-wide): a ranked list of **1–4**, best-first. - Each needs a non-empty `pattern` and `fit`, and a boolean `simplification_pattern`. Set - `simplification_pattern: true` **only** for a distinctive capability that collapses a whole legacy - pattern (a managed connector, declarative `AUTO CDC`, Auto Loader, system tables replacing a - home-grown logging tier) — not for a like-for-like port. Rank the `true` patterns first. +- **`recommended_patterns`** (per pipeline and system-wide): **1–4** patterns. The validator enforces + the count and that each has a non-empty `pattern` and `fit` and a boolean `simplification_pattern`; + it does **not** enforce ordering. Set `simplification_pattern: true` **only** for a distinctive + capability that collapses a whole legacy pattern (a managed connector, declarative `AUTO CDC`, Auto + Loader, system tables replacing a home-grown logging tier) — not for a like-for-like port. By + convention (not validated), order them best-first and list the `simplification_pattern: true` ones + ahead of like-for-like ports. - **`system_recommendation`** needs a non-empty `headline` and a `recommended_patterns` list; `cascade` (non-empty strings) and `decision_driver` are optional. - **Relationship edges** come in two tiers: diff --git a/skills/flowx-migrate/SKILL.md b/skills/flowx-migrate/SKILL.md index e12b05c..e90467e 100644 --- a/skills/flowx-migrate/SKILL.md +++ b/skills/flowx-migrate/SKILL.md @@ -2,7 +2,8 @@ name: flowx-migrate description: > End-to-end migration of a source orchestrator's pipelines (Azure Data Factory, Apache Airflow) - to Databricks Lakeflow Jobs. Orchestrates discover, convert, and package phases in sequence. + to Databricks Lakeflow Jobs. Orchestrates discover → enrich → route → convert/fill → package in + sequence. triggers: - "migrate pipelines" - "migrate ADF" @@ -16,16 +17,25 @@ triggers: # End-to-End Source to Databricks Migration Orchestrate the complete migration of a source orchestrator's pipelines to Databricks Lakeflow Jobs -via Declarative Automation Bundles. This skill runs all three phases in sequence: discover, convert, -package. +via Declarative Automation Bundles. This skill runs the full flow in sequence: discover → enrich +(default) → route → convert/fill → package. ## Context This is the top-level orchestration skill. It runs the full migration pipeline: 1. **Discover** — Parse the source's definitions into a typed inventory -2. **Convert** — Convert the source's tasks to Databricks IR (deterministic + agentic) -3. **Package** — Generate Databricks Declarative Automation Bundles for deployment +2. **Enrich** *(default)* — Author the agentic `insights` layer over the inventory and merge it + (`flowx-enrich`); skippable for a deterministic-only pass +3. **Convert (deterministic baseline)** — Convert the source's tasks to Databricks IR, producing + `.work/translation_report.json`. This runs **once, before routing** (routing reads and edits it); + route can also trigger it in-process (`flowx-convert`) +4. **Route + fill** — Decide, per connected component, deterministic vs. agentic conversion; record + the fingerprint-bound `metadata/conversion_plan.json`; edit the baseline report so routed-agentic + groups become placeholder gaps; then **fill** those gaps additively — per-pipeline + `convert --merge-agentic` or cross-pipeline `fill-agentic combine`. **Never re-run a plain + `convert` after routing** — it would rewrite the report and erase the placeholders (`flowx-route`) +5. **Package** — Generate Databricks Declarative Automation Bundles for deployment Each phase builds on the output of the previous phase. The user is shown a summary and asked to confirm before proceeding to the next phase. @@ -40,7 +50,7 @@ source path (`--adf-source-path` / `--airflow-source-path`, both aliases of `--s ## How to run this skill — MCP tools or venv CLI -This skill orchestrates all three phases. Run the **`setup`** skill first if you haven't. There are +This skill orchestrates the full flow. Run the **`setup`** skill first if you haven't. There are two execution paths: ### MCP tools (Databricks Genie Code, or a local stdio registration) @@ -101,8 +111,11 @@ discover/convert and for `inputs discover`/`inputs convert`; for Airflow, swap ` ``` flowx(command="inputs", parameters={"phase": "discover", "source": "adf"}) # source req for discover/convert flowx(command="discover", parameters={"source": "adf", "adf_definitions": {...}, "output_dir": ..., "pipeline": ...}) +flowx(command="enrich", parameters={"output_dir": ..., "insights": {...}}) # default: author + merge the insights layer (see flowx-enrich) +flowx(command="route", parameters={"output_dir": ...}) # recommend; re-call with "plan": {...} to record + edit the report (see flowx-route) flowx(command="convert", parameters={"source": "adf", "output_dir": ..., "pipeline": ...}) -flowx(command="merge_agentic", parameters={"source": "adf", "report_path": ..., "agentic_results_dir": ..., "output_path": ...}) # ADF only, if agentic results +flowx(command="merge_agentic", parameters={"source": "adf", "report_path": ..., "agentic_results_dir": ..., "output_path": ...}) # ADF only, per-pipeline agentic fill +flowx(command="fill_agentic", parameters={"output_dir": ..., "members": [...], "pipelines": [...]}) # cross-pipeline combine of a routed-agentic group flowx(command="inspect", parameters={"report_path": ...}) flowx(command="apply_answers", parameters={"report_path": ..., "answers": [...], "output_dir": ...}) flowx(command="package", parameters={"output_dir": ..., "catalog": ..., "schema": ...}) @@ -209,15 +222,60 @@ If the user says no, explain the options: - Review `/metadata/inventory.json` (and `profile_report.csv`) to understand unsupported activities and pipeline complexity - Manually classify activities before proceeding -If the user says yes, proceed to step 4. - -### Step 4 — Phase 2: Convert - -Invoke the `flowx:flowx-convert` skill with: +If the user says yes, proceed to step 3.5. + +### Step 3.5 — Enrich the inventory (default) + +By default, chain into enrichment: invoke the **`flowx:flowx-enrich`** skill to author the agentic +`insights` layer (factory-wide recommendation, per-pipeline intent + recommended Databricks patterns, +cross-pipeline relationships) and merge it into `/metadata/inventory.json`. flowx has no +LLM — you author the insights JSON and `enrich` validates + merges it additively, leaving every +deterministic inventory key byte-identical. The routing step consumes this block to present the +agentic conversion option per group. + +**Skip only for a deterministic-only pass.** When the user explicitly asked for a headless, +deterministic-only migration (no agent/LLM authoring), skip enrich — the deterministic +`inventory.json` is complete and `flowx-route` still works from the structure alone. Otherwise enrich +by default. + +### Step 3.6 — Route the conversion (deterministic vs. agentic) + +Routing reads and edits the **deterministic baseline** report `.work/translation_report.json`, so +convert (Step 4) must have produced it first. The clean path is to let `route` **trigger convert once +in-process** — pass `--source` / `--source-path` and route builds the baseline and then edits it in a +single step (this *is* the Phase-2 convert; a separate Step 4 run is then unnecessary). Otherwise run +Step 4 before this step. + +Invoke the **`flowx:flowx-route`** skill to present the per-connected-component recommendation, take +the **customer's** decision (interactive, or an authored plan — the customer decides +deterministic-vs-agentic per component; the agent presents the options and serializes only the +approved decision), and record the fingerprint-bound `/metadata/conversion_plan.json`. +Recording a plan edits `.work/translation_report.json` so every routed-**agentic** group's tasks +become placeholder gaps; deterministic groups (and a no-agentic-route plan) leave the report +untouched — non-breaking. + +**Once routing has recorded agentic placeholders, never run a plain `convert` again** — a second +`convert --source-dir` overwrites `.work/translation_report.json` and erases the placeholders. +Routed-agentic groups are **filled** additively in Step 5.05 (per-pipeline `convert --merge-agentic` +or cross-pipeline `fill-agentic combine`), which is the only convert after routing. If the user wants +a straight deterministic migration, they accept the all-deterministic recommendation here and the +rest of the flow is unchanged. + +### Step 4 — Phase 2: Convert (deterministic baseline, runs before routing) + +Convert produces the deterministic baseline report `/.work/translation_report.json` that +routing (Step 3.6) reads and edits — so it runs **before** the plan is recorded. If you let `route` +trigger convert in-process (Step 3.6), this step is already done; otherwise invoke the +`flowx:flowx-convert` skill directly, before routing, with: - `--source `: the same source discover used - Source path: the original source path (same one discover used) - Output dir: the same shared `` (convert writes its report to `/.work/`) +> **Run this convert exactly once, before routing.** After Step 3.6 records agentic placeholders, do +> **not** re-run a plain `convert` — it rewrites `.work/translation_report.json` and erases the +> routed-agentic placeholders. The only convert after routing is the additive `convert --merge-agentic` +> fill in Step 5.05. + Wait for the translation to complete and present the summary: ``` @@ -241,6 +299,25 @@ For failures, suggest: - Retry with additional context - Skip and add placeholder +### Step 5.05 — Fill routed-agentic groups (before just-in-time config) + +If Step 3.6 routed any component **agentic**, its pipelines' tasks are now `PlaceholderActivity` +nodes with one tagged gap each. Fill them via the **`flowx:flowx-route`** skill (Step 3 there) before +the just-in-time config below, so the filled tasks get configuration-stamped and packaged: + +- **Per-pipeline** (the pipeline stays 1:1): author one result JSON per gap and merge with + `convert --source adf --merge-agentic --report .work/translation_report.json --agentic-results ` + (ADF only; `--source` is mandatory for the convert phase — it exits 2 without it). This merge is + additive: it replaces only the placeholder tasks and leaves every other pipeline byte-identical. + Airflow per-gap fills use the `flowx-resolve-airflow-gaps` skill instead. +- **Cross-pipeline COMBINE** (N pipelines → M, e.g. one Lakeflow Connect pipeline): author the + replacement pipeline IR (typically with `AgenticComponentActivity` nodes) and run + `fill-agentic combine --output-dir --members "" --pipelines-path `; + the members must exactly match the routed-agentic component and the merged report is validated + structurally before it is written. + +Skip this step entirely when no component was routed agentic. + ### Step 5.1 — Gather just-in-time translation configuration Run `inspect` **once** to get the full option schema (every option carries a `show_when` condition), @@ -432,7 +509,7 @@ See `references/workflow.md` for a detailed description of the three-phase archi ## Output Artifacts -All three phases write into a single shared `` (default `./flowx_output`): +All phases write into a single shared `` (default `./flowx_output`): | Path | Phase | Contents | |---|---|---| diff --git a/skills/flowx-migrate/references/workflow.md b/skills/flowx-migrate/references/workflow.md index faf083b..cb345ef 100644 --- a/skills/flowx-migrate/references/workflow.md +++ b/skills/flowx-migrate/references/workflow.md @@ -60,6 +60,28 @@ ADF JSON Exports - Datasets and linked services are parsed for context but not independently translated — they inform the activity translators. - Triggers are included in the inventory and translated in phase 2. +## Between Discover and Convert: Enrich (default) + Route + +**Skills:** `flowx:flowx-enrich`, `flowx:flowx-route` + +After discover writes the deterministic inventory, the standard flow enriches and routes before +convert. Both are additive and contain **no LLM** — the agent authors, the library validates and +merges: + +- **Enrich (default):** the agent authors an `insights` layer (factory-wide recommendation, + per-pipeline intent + recommended Databricks patterns, cross-pipeline relationships) and `enrich` + merges it into `metadata/inventory.json` under a single additive `insights` key, leaving every + deterministic key byte-identical. Skippable for a deterministic-only, headless pass. +- **Route:** groups pipelines into connected components over control lineage and, per component, + records a `deterministic` or `agentic` decision as the fingerprint-bound + `metadata/conversion_plan.json`. Recording an agentic decision edits `.work/translation_report.json` + so those groups' tasks become placeholder gaps; a fully-deterministic plan (or no plan) leaves + convert/package behaving exactly as before — the non-breaking guarantee. + +Routed-agentic groups are then filled during convert: per-pipeline via `convert --merge-agentic`, or +cross-pipeline (N→M, e.g. a Lakeflow Connect collapse) via `fill-agentic combine` using authored +`AgenticComponentActivity` nodes. + ## Phase 2: Convert **Skill:** `flowx:flowx-convert` diff --git a/skills/flowx-route/SKILL.md b/skills/flowx-route/SKILL.md new file mode 100644 index 0000000..f1e2ad3 --- /dev/null +++ b/skills/flowx-route/SKILL.md @@ -0,0 +1,262 @@ +--- +name: flowx-route +description: > + Route each connected component of the discovered inventory to a deterministic (1:1 engine) or + agentic (LLM-assisted re-architecture) conversion, record the fingerprint-bound conversion plan, + and fill the routed-agentic groups — per-pipeline via convert --merge-agentic, or cross-pipeline + via fill-agentic combine. Runs after enrich, before/with convert. +triggers: + - "route pipelines" + - "route conversion" + - "conversion plan" + - "deterministic or agentic" + - "fill agentic" + - "combine pipelines" + - "agentic conversion" + - "recommend conversion route" +--- + +# Route the Conversion (deterministic vs. agentic) and Fill the Gaps + +After discover (and, by default, `flowx-enrich`), routing decides **per connected component** whether +each part of the factory converts **deterministically** (the typed engine's 1:1 translation) or +**agentically** (an LLM-authored re-architecture, e.g. collapsing five extractors onto one Lakeflow +Connect pipeline). It records the decision as a fingerprint-bound `metadata/conversion_plan.json`, +edits the translation report so routed-agentic groups become placeholder gaps, and then you author +the fill. + +**There is no LLM inside flowx.** The library computes the recommendation deterministically and only +*validates and records* the decision and the authored fill — the same author → validate → merge +contract `enrich` uses. This is additive and non-breaking: with **no recorded plan**, `convert` and +`package` behave exactly as before. + +**Who decides what.** The library **recommends** a route per component; the **customer decides** +deterministic-vs-agentic for each component; the **agent** presents the options and the +recommendation, serializes only the customer's *approved* decision into the plan, and authors the +agentic fill. The agent never picks the route on the customer's behalf — it records the customer's +choice and does the mechanical work of the approved fill. + +## Where routing sits + +``` +discover → enrich (default) → convert (deterministic baseline) → route (decide + edit) → fill-agentic → package +``` + +Convert builds the deterministic baseline report **before** routing (route can trigger it in-process +via `--source` / `--source-path`); routing then edits that report. After routing has recorded agentic +placeholders, the only convert is the additive `convert --merge-agentic` fill — never a second plain +`convert`, which would overwrite the report and erase the placeholders. + +Pipelines are grouped into weak/undirected **connected components** over the inventory's control +lineage (`lineage.control_edges`), so mutually-referencing pipelines are decided together and a +caller/callee reference is never split across incompatible routes. Both ADF (`ExecutePipeline`) and +Airflow (`RunJob`) emit control edges, so routing is source-neutral. + +Run the **`setup`** skill first if you haven't. Everything below has an MCP-tool path (Genie Code, or +a local stdio registration — call the single **`flowx`** tool, run no `python3`/`$PY`) and a venv-CLI +path (local; `PY="$(cat /.migration-venv)"` and `export PYTHONPATH="/src"`). + +## Step 1 — Recommend (the dry run you read first) + +Call `route` with **no decision** to get the recommendation: every component with its `members`, both +first-class conversion `options` (a *deterministic* option carrying the engine-capability assessment ++ any uncovered gaps, and an *agentic* option carrying the `recommended_patterns` from `enrich`, with +any `simplification_pattern` flagged), the `findings` (unresolved/dangling control edges kept, never +severed), and a ready-to-record `default_plan` proposing `decision == recommended` for every +component. + +- **MCP tool:** `flowx(command="route", parameters={"output_dir": ""})` +- **venv CLI:** + + ```bash + export PYTHONPATH="/src" + PY="$(cat /.migration-venv)" + "$PY" -m flowx.adapter route --output-dir + ``` + + On a non-TTY with no `--plan-path`, this emits the recommendation and exits 0 (it edits nothing). + `route` reads `metadata/inventory.json`; run discover first. `--out ` writes the JSON to a + file instead of stdout. + +`recommended` is the library's starting suggestion — `deterministic` when the whole component is +engine-capable, else `agentic`. Present each component's members, its recommendation, and the agentic +option's patterns (flag any `has_simplification` re-architecture prominently). + +## Step 2 — Decide and record the plan + +Take the per-component decision and record it. `route` validates the plan, writes the fingerprint-bound +`metadata/conversion_plan.json`, and then edits `.work/translation_report.json` + `gaps.json` so every +routed-**agentic** group's tasks become placeholder gaps (deterministic groups are left byte-identical; +when nothing is routed agentic the report is untouched — the non-breaking guarantee). + +There are three ways to supply the decision: + +- **Interactive prompt (TTY).** Run `route` with no `--plan-path` on a real terminal and it asks, per + component, `route [d]eterministic / [a]gentic (default=)`. An empty answer accepts the + recommendation. This is the from-the-seat path. +- **Authored plan file** — `--plan-path ` (or `--plan-path -` to read the plan JSON from stdin). +- **MCP tool** — pass the plan inline as `plan` (or `plan_path`): + + ``` + flowx(command="route", parameters={"output_dir": "", "plan": { ...authored plan... }}) + ``` + +### The plan shape + +The plan carries **only** the decision (and an optional rationale) per component — and the decision is +the **customer's**, not yours. Your job is to present each component's members, its `recommended` +route, and both `options`, then serialize the customer's pick; you do not choose the route yourself. +The library recomputes `members`, `recommended`, and both `options` on record, so recorded facts +cannot drift from the inventory or be faked. Start from the recommendation's `default_plan` and flip +the components the customer chose to override: + +```json +{ + "components": [ + {"component_id": "component-1", "members": ["IngestSalesforce", "IngestWorkday"], "decision": "agentic", + "rationale": "Collapse both extractors onto one Lakeflow Connect pipeline"}, + {"component_id": "component-2", "members": ["BuildMart"], "decision": "deterministic"} + ] +} +``` + +Rules the validator enforces (all violations returned at once; nothing written on failure): + +- `decision` must be `"deterministic"` or `"agentic"`; `rationale` (optional) must be a non-empty + string when present. +- Every `member` must be a real inventory pipeline, and a component's `members` must **exactly match + one computed connected component** — a decision can never split a component or span two. +- The plan is a **bijection**: every component is decided exactly once (no partial plan, no + duplicate/conflicting decisions). +- `component_id` is **required** on every component — a non-empty string that must match the computed + component for those `members`. Start from the recommendation's `default_plan`, which already carries + the correct `component_id` for each component. + +### Triggering convert if the report is missing + +`route` edits `.work/translation_report.json`, which the convert phase produces. If it is missing, +pass `--source ` **and** `--source-path ` (MCP: `source` plus the source-specific +path param — `adf_source_path` or `airflow_source_path`; a generic `source_path` is **ignored**, and +ADF also accepts `adf_definitions` / `adf_volume_path` / `adf_workspace_path`) and `route` triggers +the convert phase in-process first. Otherwise run `flowx-convert` before routing. + +### venv CLI + +```bash +"$PY" -m flowx.adapter route --output-dir --plan-path plan.json \ + [--source adf --source-path ] # only needed to trigger convert when the report is missing +``` + +Exit 1 (nothing written) on a missing report it cannot produce, or a plan that fails validation. + +## Step 3 — Fill the routed-agentic groups + +Every routed-agentic pipeline's tasks are now `PlaceholderActivity` nodes with one pipeline-tagged +`AgenticGap` each. Two fills exist depending on the grain; you author the replacement (no LLM in the +library — it validates and merges). + +### 3a — Per-pipeline fill (keep the pipeline 1:1) + +When a routed-agentic pipeline stays one pipeline and you just author a Databricks task per gap, reuse +the **existing** name-matched merge — the same path the agentic gap flow uses. Write one JSON file per +resolved gap into an agentic-results directory: + +```json +{ + "activity_name": "", + "pipeline": "", + "task": {"type": "NotebookActivity", "name": "", "task_key": "", + "notebook_path": "/Workspace/.../your_translated_notebook"} +} +``` + +Then merge (`task_key`/`depends_on` are inherited from the placeholder when omitted, preserving edges): + +```bash +"$PY" -m flowx.adapter convert --source adf --merge-agentic \ + --report /.work/translation_report.json \ + --agentic-results \ + [--output ] # default: overwrite --report +``` + +`--source` is **mandatory** for the convert phase — the phase runner exits 2 without it, even for +`--merge-agentic`. Use `--source adf` here (`merge_agentic` is **ADF-only**; an Airflow per-gap fill +would use `--source airflow`, but Airflow gaps normally go through the `flowx-resolve-airflow-gaps` +skill). This merge is **additive** — it replaces only the placeholder tasks in the existing report and +leaves every other pipeline byte-identical, so it never erases routing's edits. + +MCP: `flowx(command="merge_agentic", parameters={"source": "adf", "report_path": ..., "agentic_results_dir": ..., "output_path": ...})`. +Placeholders are replaced in place, status → `translated`; exits non-zero if any result can't be +matched. + +### 3b — Cross-pipeline COMBINE (N pipelines → M) + +When the decision is a re-architecture that changes pipeline count — e.g. five extractor pipelines +collapse onto **one** Lakeflow Connect pipeline — use `fill-agentic combine`. You author the +replacement pipeline(s) as IR dicts and the whole routed group is swapped for them. + +```bash +"$PY" -m flowx.adapter fill-agentic combine \ + --output-dir \ + --members "IngestSalesforce,IngestWorkday" \ + --pipelines-path authored_pipelines.json \ + [--out ] +``` + +MCP: `flowx(command="fill_agentic", parameters={"output_dir": ..., "members": [...], "pipelines": [...]})` +(pass `pipelines` inline as a list, or `pipelines_path`). + +`combine` is the only action. Its guarantees: + +- `--members` (comma-separated; MCP accepts a list) must **exactly match** a routed-**agentic** + component in the recorded `metadata/conversion_plan.json`, whose `inventory_sha256` must still match + the current inventory. A partial group, a superset, a typo, or a deterministic component is refused + — you can't swap pipelines the plan didn't route agentic. +- `--pipelines-path` is a JSON **list** of pipeline IR dicts (the authored replacements), typically + carrying `AgenticComponentActivity` nodes (see below). +- The merged report is **always** validated with the structural bundle invariants (a real + `prepare → write_bundle` pass) before it is written — no bypass — so a duplicate key, dangling + dependency, cycle, or dangling pipeline/run_job reference can never land on disk. On any violation, + `ok` is `false`, `violations` lists them, and nothing is written. + +### Authoring an `AgenticComponentActivity` (the escape hatch) + +When the target can't be expressed by the typed engine (e.g. a managed Lakeflow Connect ingestion +pipeline), emit an `AgenticComponentActivity` task inside the authored pipeline. It carries the raw +bundle components the package phase writes verbatim: + +```json +{ + "name": "IngestAll", + "task_key": "ingest_all", + "type": "AgenticComponentActivity", + "files": [ + {"path": "src/ingest/lakeflow_connect.py", "content": "# authored pipeline source ..."} + ], + "resources": [ + {"resource_key": "ingest_all_pipeline", "definition": { ...raw pipeline resource... }} + ], + "task": {"pipeline_task": {"pipeline_id": "${resources.pipelines.ingest_all_pipeline.id}"}}, + "raw_definition": { ...original source definitions, retained for auditing... } +} +``` + +- `files` — files written below the bundle `src/`; each is `path` + either UTF-8 `content` or + base64 `binary_content`. +- `resources` — bundle resources in the `resource_key` + raw `definition` shape the bundle writer + expects. +- `task` — the raw Databricks task fragment (`pipeline_task` or `notebook_task`) wiring to an + authored resource or file. +- `raw_definition` — the original source definition, retained for provenance. + +## Step 4 — Continue to package + +Once the routed-agentic groups are filled and the report validates, continue with `flowx-convert`'s +just-in-time configuration (`inspect`/`modify`) as usual, then `flowx-package`. The recorded +`metadata/conversion_plan.json` is kept alongside `inventory.json` as the routing record. + +## Reference + +- `flowx-enrich` skill — the `insights` layer routing's agentic option consumes. +- `flowx-convert` — the deterministic engine, per-pipeline `merge_agentic`, and just-in-time config. +- `flowx-package` — turns the (filled) report into the deployable DAB bundle.