Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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: 2 additions & 0 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,8 @@ intermediates under `.work/` (pruned by `package`).
| `reporting/coverage.py` | Builds per-pipeline coverage rows from `metadata/` |
| `reporting/results.py` | Writes per-run coverage to a UC table (run_id/run_date/run_by) via the SDK |
| `reporting/dashboard.py` | Installs + publishes an AI/BI coverage dashboard over the results table |
| `routing.py` | Groups pipelines into connected components over control lineage; recommends deterministic/agentic per component (both options) and records the user's decision as `metadata/conversion_plan.json` (additive; convert untouched) |
| `models/conversion_plan.py` | Source-neutral conversion-plan artifact model (per-component decision + both options), the routing counterpart to `models/insights.py` |

## Activity Types

Expand Down
79 changes: 76 additions & 3 deletions src/flowx/adapter/__main__.py
Original file line number Diff line number Diff line change
@@ -1,9 +1,9 @@
"""Unified CLI entry point that the flowx skills and MCP tools drive via subprocesses.

Exposes stateless subcommands -- the ``discover``/``convert``/``package`` phase runners plus
``inspect``, ``modify``, ``resolve-agentic``, ``enrich``, ``inputs``, ``materialize-lookup``,
``workspace-paths``, ``record-results``, and ``install-dashboard`` -- so each agent turn runs as an
independent process holding no session state across user prompts.
``inspect``, ``modify``, ``resolve-agentic``, ``enrich``, ``route``, ``inputs``,
``materialize-lookup``, ``workspace-paths``, ``record-results``, and ``install-dashboard`` -- so each
agent turn runs as an independent process holding no session state across user prompts.
"""

from __future__ import annotations
Expand Down Expand Up @@ -87,6 +87,8 @@ def main(argv: list[str] | None = None) -> int:
return _run_resolve_agentic(args)
if args.command == "enrich":
return _run_enrich(args)
if args.command == "route":
return _run_route(args)
if args.command == "record-results":
return _run_record_results(args)
if args.command == "install-dashboard":
Expand Down Expand Up @@ -163,6 +165,46 @@ def _run_enrich(args: argparse.Namespace) -> int:
return 0 if result.get("ok") else 1


def _run_route(args: argparse.Namespace) -> int:
"""Implements ``route``: compute a per-connected-component conversion recommendation (``recommend``)
or validate and record an agent-authored conversion plan (``record``).

``recommend`` is read-only: it emits the components, both conversion options per component, the
findings, and a ready-to-record default plan to stdout. ``record`` validates the authored plan and,
only when clean, writes ``metadata/conversion_plan.json``; it emits the result JSON (``ok`` /
``violations`` / counts) and returns 1 on validation failure (plan left untouched).
"""
from flowx import routing

inventory_path = args.output_dir / "metadata" / "inventory.json"
if not inventory_path.exists():
print(f"No inventory.json under {inventory_path.parent}; run the discover phase first.", file=sys.stderr)
return 1

if args.action == "recommend":
try:
inventory = json.loads(inventory_path.read_text(encoding="utf-8"))
except (OSError, json.JSONDecodeError) as error:
print(f"Failed to read {inventory_path}: {error}", file=sys.stderr)
return 1
_emit_json(routing.build_recommendation(inventory), args.out)
return 0

if args.plan_path is None:
print("route record requires --plan-path.", file=sys.stderr)
return 2
try:
result = routing.record_plan(args.output_dir, plan_path=args.plan_path)
except FileNotFoundError as error:
print(str(error), file=sys.stderr)
return 1
except (OSError, ValueError, json.JSONDecodeError) as error:
print(f"Failed to record conversion plan: {error}", file=sys.stderr)
return 1
_emit_json(result, args.out)
return 0 if result.get("ok") else 1


def _run_record_results(args: argparse.Namespace) -> int:
"""Implements ``record-results``: write per-pipeline coverage to a UC table.

Expand Down Expand Up @@ -521,6 +563,37 @@ def _build_parser() -> argparse.ArgumentParser:
help="Optional output file for the enrich result JSON; defaults to stdout.",
)

route = subparsers.add_parser(
"route",
help=(
"Recommend a per-connected-component deterministic/agentic conversion route, or record "
"the user's decision as metadata/conversion_plan.json (additive; convert is untouched)."
),
)
route.add_argument(
"action",
choices=("recommend", "record"),
help="'recommend' emits components + both options + a default plan; 'record' validates and writes the plan.",
)
route.add_argument(
"--output-dir",
type=Path,
required=True,
help="Migration output directory (reads metadata/inventory.json; record writes metadata/conversion_plan.json).",
)
route.add_argument(
"--plan-path",
type=Path,
default=None,
help="For 'record': path to the agent-authored conversion plan JSON to validate and record.",
)
route.add_argument(
"--out",
type=Path,
default=None,
help="Optional output file for the recommendation / result JSON; defaults to stdout.",
)

record = subparsers.add_parser(
"record-results",
help="Write per-pipeline migration coverage for this run to a Unity Catalog table.",
Expand Down
45 changes: 45 additions & 0 deletions src/flowx/mcp/server.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@ def _transport_security() -> TransportSecuritySettings:
flowx("inputs", {"phase": "discover", "source": "adf"}) # learn a phase's inputs
flowx("discover", {"source": "adf", "adf_source_path": "...", "output_dir": "..."})
flowx("enrich", {"output_dir": "...", "insights": {...}}) # optional: merge agent-authored insights
flowx("route", {"output_dir": "...", "action": "recommend"}) # optional: per-component route + user decision
flowx("convert", {"source": "adf", "output_dir": "..."})
flowx("inspect", {"report_path": "<output_dir>/.work/translation_report.json"})
flowx("apply_answers", {"report_path": "...", "answers": ["id=value"], "output_dir": "..."})
Expand Down Expand Up @@ -326,6 +327,42 @@ def _run(path: str) -> dict[str, Any]:
return _run(str(inline_path))


def _cmd_route(p: dict[str, Any]) -> dict[str, Any]:
"""Recommend a per-connected-component conversion route, or record the user's decision.

``action`` (default ``"recommend"``): ``recommend`` emits the components, both conversion options
per component, findings, and a ready-to-record default plan. ``record`` validates an authored plan
and writes ``metadata/conversion_plan.json``; supply the plan either inline as ``plan`` (a JSON
object) or via ``plan_path`` (exactly one, staged to a temp file so the same CLI contract runs on
both paths). The returned ``ok`` reflects the CLI's success; ``result`` carries its JSON.
"""
output_dir = p.get("output_dir", "./flowx_output")
action = p.get("action", "recommend")
if action == "recommend":
result = runner.run_adapter(["route", "recommend", "--output-dir", output_dir])
return {"ok": result.ok, "result": runner.parse_stdout_json(result), "process": result.as_dict()}
if action != "record":
return {"ok": False, "error": f"unknown route action {action!r}; expected 'recommend' or 'record'."}

plan = p.get("plan")
plan_path = p.get("plan_path")
if (plan is None) == (plan_path is None):
return {"ok": False, "error": "provide exactly one of 'plan' (inline object) or 'plan_path'."}

def _run(path: str) -> dict[str, Any]:
result = runner.run_adapter(["route", "record", "--output-dir", output_dir, "--plan-path", path])
payload = runner.parse_stdout_json(result)
ok = bool(isinstance(payload, dict) and payload.get("ok"))
return {"ok": ok, "result": payload, "process": result.as_dict()}

if plan_path is not None:
return _run(str(plan_path))
with tempfile.TemporaryDirectory(prefix="flowx-plan-") as temporary:
inline_path = Path(temporary) / "conversion_plan.json"
inline_path.write_text(json.dumps(plan, indent=2), encoding="utf-8")
return _run(str(inline_path))


def _cmd_inspect(p: dict[str, Any]) -> dict[str, Any]:
args: list[Any] = ["inspect", p["report_path"]]
for answer in p.get("answers") or []:
Expand Down Expand Up @@ -529,6 +566,7 @@ def _cmd_install_dashboard(p: dict[str, Any]) -> dict[str, Any]:
"merge_agentic": _cmd_merge_agentic,
"resolve_agentic": _cmd_resolve_agentic,
"enrich": _cmd_enrich,
"route": _cmd_route,
"inspect": _cmd_inspect,
"apply_answers": _cmd_apply_answers,
"materialize_lookup": _cmd_materialize_lookup,
Expand Down Expand Up @@ -587,6 +625,13 @@ def flowx(command: str, parameters: dict[str, Any] | None = None) -> dict[str, A
`insights` key (atomic, idempotent). `ok` reflects validation; `result.violations` lists any
problems and the inventory is left untouched on failure. Author the insights by reading
inventory.json + the source artifacts first (see the flowx-discover skill's insights guide).
- "route": output_dir(req), action("recommend" default | "record"), one of plan(inline object)
| plan_path(for record) — group pipelines into connected components over control lineage and
route each. `recommend` emits the components with BOTH conversion options (deterministic
capability + motif/coverage evidence, and the agentic recommended patterns with any
simplification pattern surfaced), plus a ready-to-record default plan. `record` validates the
user's per-component decision and writes metadata/conversion_plan.json (additive; convert is
untouched, and with no plan behavior is exactly as today).
- "inspect": report_path(req) — return the full translation-option schema (every option with
a `show_when` condition) for the agent to walk locally. See "Collecting options" below.
- "apply_answers": report_path(req), answers(req, list of "ID=VALUE"), output_dir, lookup_csv.
Expand Down
135 changes: 135 additions & 0 deletions src/flowx/models/conversion_plan.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,135 @@
"""Conversion-plan artifact (routing, #77) -- the user-approved per-component conversion decision.

The routing step (:mod:`flowx.routing`) groups pipelines into connected components over control
lineage, presents each component's two conversion options as first-class peers, and records the
user's per-component choice as ``metadata/conversion_plan.json``. These models are **source-neutral**
and document the shape of that artifact; the validate/record engine works on the raw dict form and
these dataclasses back the unit tests, mirroring the split in :mod:`flowx.models.insights`.

The agent authors **only** :attr:`ComponentPlan.decision` (and an optional
:attr:`ComponentPlan.rationale`). Everything else -- ``component_id``, ``members``, ``recommended``,
and both :class:`ComponentOptions` -- is recomputed by the library on record so the recorded facts
can never drift from the inventory or be faked. The library also owns :attr:`ConversionPlan.schema_version`
and :attr:`ConversionPlan.inventory_sha256` (the fingerprint that binds the plan to the inventory).

This is Phase-1, descriptive-only routing metadata: recording a plan does not alter ``convert``.
"""

from __future__ import annotations

from dataclasses import dataclass, field
from typing import Any

# The plan schema version stamped onto the recorded artifact. Bump on any backwards-incompatible
# change to the recorded shape.
SCHEMA_VERSION = "1"

# The two conversion routes a component can take.
DECISION_DETERMINISTIC = "deterministic"
DECISION_AGENTIC = "agentic"
DECISIONS: tuple[str, ...] = (DECISION_DETERMINISTIC, DECISION_AGENTIC)


@dataclass(slots=True, kw_only=True)
class DeterministicOption:
"""The deterministic (1:1 engine) conversion option for a component.

Attributes:
capable: ``True`` when every activity in the component is engine-capable -- each either has a
``"deterministic"`` strategy or is claimed by a detected motif -- so the whole component
can convert deterministically with no gap.
activity_counts: Count of activities by strategy bucket (``deterministic`` / ``agentic`` /
``unsupported``) across the component's pipelines.
motifs: The motif ids detected in the component (the #64 multi-activity capability signal),
sorted and de-duplicated.
uncovered: One entry per activity that keeps the component from being fully deterministic --
a dict of ``pipeline`` / ``activity`` / ``type`` / ``strategy``. Empty when ``capable``.
"""

capable: bool
activity_counts: dict[str, int] = field(default_factory=dict)
motifs: list[str] = field(default_factory=list)
uncovered: list[dict[str, Any]] = field(default_factory=list)


@dataclass(slots=True, kw_only=True)
class AgenticPattern:
"""One recommended Databricks pattern for the agentic option, drawn from a pipeline's insights.

Attributes:
pipeline: The member pipeline the pattern was recommended for.
pattern: The named, publicly-documented Databricks capability (verbatim from the insight).
fit: One line on why it fits / what it replaces.
simplification_pattern: ``True`` when the pattern uses a distinctive capability that collapses
a whole legacy pattern (e.g. a multi-pipeline -> Lakeflow Connect re-architecture).
"""

pipeline: str
pattern: str
fit: str
simplification_pattern: bool


@dataclass(slots=True, kw_only=True)
class AgenticOption:
"""The agentic conversion option for a component.

Attributes:
recommended_patterns: The recommended patterns gathered from the member pipelines' insights,
each tagged with its pipeline. Empty when the inventory carries no insights.
has_simplification: ``True`` when any recommended pattern is a ``simplification_pattern`` --
surfaced prominently so the user sees a re-architecture option, not a buried sub-key.
"""

recommended_patterns: list[AgenticPattern] = field(default_factory=list)
has_simplification: bool = False


@dataclass(slots=True, kw_only=True)
class ComponentOptions:
"""Both conversion options for a component, as first-class peers."""

deterministic: DeterministicOption
agentic: AgenticOption


@dataclass(slots=True, kw_only=True)
class ComponentPlan:
"""One component's routing decision.

Attributes:
component_id: Stable id assigned by the library (``"component-<n>"``).
members: The component's pipeline names, sorted (library-computed).
recommended: The library's starting suggestion -- ``"deterministic"`` when the component is
engine-capable, else ``"agentic"``.
decision: The user's authored per-component choice (may override :attr:`recommended`).
options: Both conversion options with their evidence (library-computed). Optional here so the
authored input -- which carries only the decision -- can round-trip through this model.
rationale: Optional author note on why this decision was chosen.
"""

component_id: str
members: list[str] = field(default_factory=list)
recommended: str | None = None
decision: str
options: ComponentOptions | None = None
rationale: str | None = None


@dataclass(slots=True, kw_only=True)
class ConversionPlan:
"""The recorded conversion-plan artifact.

Attributes:
components: One :class:`ComponentPlan` per connected component.
findings: Human-readable notes about unresolved/dangling control edges retained during
component computation (never silently severed).
schema_version: Library-owned plan schema version.
inventory_sha256: Library-owned fingerprint binding the plan to the deterministic inventory
base (see :func:`flowx.discovery_insights.inventory_fingerprint`).
"""

components: list[ComponentPlan] = field(default_factory=list)
findings: list[str] = field(default_factory=list)
schema_version: str = SCHEMA_VERSION
inventory_sha256: str | None = None
Loading