diff --git a/AGENTS.md b/AGENTS.md index e1ae8c4..7f2da25 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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 diff --git a/src/flowx/adapter/__main__.py b/src/flowx/adapter/__main__.py index a2bb308..372d781 100644 --- a/src/flowx/adapter/__main__.py +++ b/src/flowx/adapter/__main__.py @@ -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 @@ -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": @@ -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. @@ -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.", diff --git a/src/flowx/mcp/server.py b/src/flowx/mcp/server.py index 0651640..79afd92 100644 --- a/src/flowx/mcp/server.py +++ b/src/flowx/mcp/server.py @@ -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": "/.work/translation_report.json"}) flowx("apply_answers", {"report_path": "...", "answers": ["id=value"], "output_dir": "..."}) @@ -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 []: @@ -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, @@ -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. diff --git a/src/flowx/models/conversion_plan.py b/src/flowx/models/conversion_plan.py new file mode 100644 index 0000000..bd0dc7a --- /dev/null +++ b/src/flowx/models/conversion_plan.py @@ -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-"``). + 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 diff --git a/src/flowx/routing.py b/src/flowx/routing.py new file mode 100644 index 0000000..a6f318c --- /dev/null +++ b/src/flowx/routing.py @@ -0,0 +1,539 @@ +"""Compute a per-connected-component conversion route over the discover inventory (#77). + +The discover phase writes a deterministic ``metadata/inventory.json``; :mod:`flowx.discovery_insights` +optionally enriches it with an additive ``insights`` block. This module turns those signals into a +**routing recommendation** and records the **user's decision** as a fingerprint-bound +``metadata/conversion_plan.json`` artifact -- a Phase-1, descriptive-only step: + +* 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; +* each component surfaces **both conversion options as first-class peers** -- a *deterministic* + option (the engine-capability assessment: every activity engine-capable via its ``strategy`` or + claimed by a detected motif, plus its motif/coverage evidence and any uncovered gaps) and an + *agentic* option (the recommended Databricks patterns from the ``insights`` block, with any + ``simplification_pattern`` such as a multi-pipeline -> Lakeflow Connect re-architecture surfaced + prominently) -- so the user can choose per group; +* ``recommended`` is a library-computed starting suggestion (deterministic when the whole component + is engine-capable), and ``decision`` is the user's authored per-component choice. + +Like :mod:`flowx.discovery_insights`, there is **no LLM here**: the tool computes the recommendation +deterministically and only *validates and records* the agent-authored decision. The agent authors +**only** the decision (and an optional rationale); the library recomputes ``members``, +``recommended``, and both options' evidence on record so they can never drift from the inventory or +be faked. + +The recorded plan is bound to the inventory via :func:`flowx.discovery_insights.inventory_fingerprint` +-- a SHA-256 over the deterministic inventory base (the ``insights`` block excluded). That base is +exactly the structural signal that decides component membership and engine capability (pipelines, +control lineage, per-activity strategy, motifs); insights are advisory evidence for the agentic +option, not part of the binding. The write is atomic (temp file + ``os.replace``) and idempotent, so +re-recording the same decision against the same inventory rewrites byte-identical bytes. + +This is additive, opt-in routing metadata only: with no recorded plan, ``convert`` and ``package`` +behave exactly as today. Component computation reads ``lineage.control_edges``, which both ADF +(``ExecutePipeline``) and Airflow (``RunJob``) emit, so the artifact is source-neutral. +""" + +from __future__ import annotations + +import json +import os +from collections import Counter +from collections.abc import Iterable +from pathlib import Path +from typing import Any + +from flowx.discovery_insights import inventory_fingerprint +from flowx.models.conversion_plan import DECISIONS, SCHEMA_VERSION + +# The strategy value that marks an activity as individually engine-capable (a 1:1 deterministic +# translation). Any other value (``"agentic"`` / ``"unsupported"`` / missing) is a gap unless the +# activity is claimed by a detected motif -- the multi-activity capability signal from #64. +_DETERMINISTIC_STRATEGY = "deterministic" + +# The recorded conversion-plan artifact lives beside inventory.json under metadata/. +PLAN_FILENAME = "conversion_plan.json" + +# Authored top-level keys (everything else the library owns and rejects on input). +_PLAN_TOP_KEYS = {"components"} +_LIBRARY_TOP_KEYS = {"schema_version", "inventory_sha256", "findings"} +# Authored per-component keys vs the fields the library recomputes and rejects on input. +_COMPONENT_AUTHORED_KEYS = {"component_id", "members", "decision", "rationale"} +_COMPONENT_LIBRARY_KEYS = {"recommended", "options"} + + +def _pipeline_names(inventory: dict[str, Any]) -> set[str]: + """The set of real pipeline names in the inventory (the foreign-key domain). + + A local copy of the same projection :mod:`flowx.discovery_insights` uses, kept here so routing + does not depend on that module's private helpers. + """ + return { + str(pipeline["name"]) + for pipeline in inventory.get("pipelines", []) + if isinstance(pipeline, dict) and pipeline.get("name") is not None + } + + +def _control_edges(inventory: dict[str, Any]) -> Iterable[tuple[str, str, str, bool]]: + """Yield ``(source_workflow, target_workflow, via_task_key, resolved)`` for every control edge. + + Lineage is placed per pipeline (one block beside each pipeline's ``activities``), so every + pipeline's ``lineage.control_edges`` are gathered into one stream, preserving inventory order. + """ + for pipeline in inventory.get("pipelines", []): + if not isinstance(pipeline, dict): + continue + lineage = pipeline.get("lineage") or {} + for edge in lineage.get("control_edges", []): + if not isinstance(edge, dict): + continue + yield ( + str(edge.get("source_workflow") or ""), + str(edge.get("target_workflow") or ""), + str(edge.get("via_task_key") or ""), + bool(edge.get("resolved", True)), + ) + + +def build_components(inventory: dict[str, Any]) -> tuple[list[list[str]], list[str]]: + """Group pipelines into weak/undirected connected components over control lineage. + + Two pipelines share a component when a **resolved** control edge joins them in either direction + and both endpoints are real inventory pipelines. Isolated pipelines each form their own + singleton component. An edge whose callee is unresolved (``resolved`` is False) or names a + pipeline absent from the inventory is **not** used to join anything and is **not** silently + severed -- it is recorded as a finding so a partial export never quietly collapses two + components into one or drops a coupling. + + Args: + inventory: The discover ``inventory.json`` document (deterministic or enriched). + + Returns: + ``(components, findings)`` where ``components`` is a list of member lists -- each sorted, + the whole list ordered by first member -- so the output is deterministic and idempotent for + a given inventory; and ``findings`` is a sorted list of human-readable notes about + unresolved/dangling edges. + """ + names = _pipeline_names(inventory) + parent: dict[str, str] = {name: name for name in names} + + def find(node: str) -> str: + root = node + while parent[root] != root: + root = parent[root] + while parent[node] != root: + parent[node], node = root, parent[node] + return root + + def union(left: str, right: str) -> None: + left_root, right_root = find(left), find(right) + if left_root != right_root: + # Attach the lexicographically larger root under the smaller for a stable shape. + low, high = sorted((left_root, right_root)) + parent[high] = low + + findings: set[str] = set() + for source, target, via, resolved in _control_edges(inventory): + if not resolved or not target or target not in names: + shown_target = target or "" + findings.add( + f"unresolved control edge {source!r} -> {shown_target!r} (via {via!r}) " + f"kept as a finding, not severed; {source!r} is routed within its own component" + ) + continue + if source in names: + union(source, target) + + groups: dict[str, list[str]] = {} + for name in names: + groups.setdefault(find(name), []).append(name) + components = sorted((sorted(members) for members in groups.values()), key=lambda members: members) + return components, sorted(findings) + + +def _activities_by_pipeline(inventory: dict[str, Any]) -> dict[str, list[dict[str, Any]]]: + """Map each pipeline name to its inventory ``activities`` list (flattened, as emitted).""" + result: dict[str, list[dict[str, Any]]] = {} + for pipeline in inventory.get("pipelines", []): + if isinstance(pipeline, dict) and pipeline.get("name") is not None: + activities = pipeline.get("activities") or [] + result[str(pipeline["name"])] = [entry for entry in activities if isinstance(entry, dict)] + return result + + +def _motifs_by_pipeline(inventory: dict[str, Any]) -> dict[str, list[dict[str, Any]]]: + """Map each pipeline name to its additive ``motifs`` list (the #64 capability signal).""" + result: dict[str, list[dict[str, Any]]] = {} + for pipeline in inventory.get("pipelines", []): + if isinstance(pipeline, dict) and pipeline.get("name") is not None: + motifs = pipeline.get("motifs") or [] + result[str(pipeline["name"])] = [entry for entry in motifs if isinstance(entry, dict)] + return result + + +def _pipeline_insights(inventory: dict[str, Any]) -> dict[str, dict[str, Any]]: + """Map each pipeline name to its ``insights.pipeline_insights`` entry, when insights are present. + + Returns an empty map when the inventory has not been enriched, so the agentic option degrades + to an empty pattern list rather than failing. + """ + insights = inventory.get("insights") + if not isinstance(insights, dict): + return {} + result: dict[str, dict[str, Any]] = {} + for entry in insights.get("pipeline_insights", []): + if isinstance(entry, dict) and entry.get("pipeline") is not None: + result[str(entry["pipeline"])] = entry + return result + + +def _deterministic_option(members: list[str], inventory: dict[str, Any]) -> dict[str, Any]: + """Assess the deterministic (1:1 engine) conversion of a component. + + An activity is engine-capable when its ``strategy`` is ``"deterministic"`` or it is claimed by a + detected motif (present in some motif's ``member_task_keys``). ``capable`` is True only when the + whole component leaves no gap; ``uncovered`` cites each activity that keeps it from being fully + deterministic. + """ + activities_by_pipeline = _activities_by_pipeline(inventory) + motifs_by_pipeline = _motifs_by_pipeline(inventory) + + motif_ids: list[str] = [] + # Motif coverage is keyed by (pipeline, activity name), never bare name: a motif claiming a task + # in one pipeline must not mark an unrelated same-named task in another pipeline of the component + # as covered, which would silently hide a real gap. + covered_activities: set[tuple[str, str]] = set() + for pipeline in members: + for motif in motifs_by_pipeline.get(pipeline, []): + if motif.get("motif_id") is not None: + motif_ids.append(str(motif["motif_id"])) + covered_activities.update((pipeline, str(key)) for key in motif.get("member_task_keys") or []) + + counts = {"deterministic": 0, "agentic": 0, "unsupported": 0} + uncovered: list[dict[str, Any]] = [] + for pipeline in members: + for activity in activities_by_pipeline.get(pipeline, []): + strategy = activity.get("strategy") + bucket = strategy if strategy in ("deterministic", "agentic") else "unsupported" + counts[bucket] += 1 + name = activity.get("name") + if strategy != _DETERMINISTIC_STRATEGY and (pipeline, name) not in covered_activities: + uncovered.append( + { + "pipeline": pipeline, + "activity": name, + "type": activity.get("type"), + "strategy": strategy, + } + ) + return { + "capable": not uncovered, + "activity_counts": counts, + "motifs": sorted(set(motif_ids)), + "uncovered": uncovered, + } + + +def _agentic_option(members: list[str], inventory: dict[str, Any]) -> dict[str, Any]: + """Surface the agent-authored recommended patterns for a component as a first-class option. + + Draws every ``recommended_patterns`` entry from the member pipelines' insights, tagging each with + its pipeline, and flags whether any is a ``simplification_pattern`` (a distinctive re-architecture + such as a multi-pipeline -> Lakeflow Connect collapse) so the user sees it prominently. + """ + insights_by_pipeline = _pipeline_insights(inventory) + recommended_patterns: list[dict[str, Any]] = [] + for pipeline in members: + insight = insights_by_pipeline.get(pipeline) + if not insight: + continue + for pattern in insight.get("recommended_patterns") or []: + if isinstance(pattern, dict): + recommended_patterns.append({"pipeline": pipeline, **pattern}) + has_simplification = any(pattern.get("simplification_pattern") for pattern in recommended_patterns) + return {"recommended_patterns": recommended_patterns, "has_simplification": has_simplification} + + +def recommend_component(members: list[str], inventory: dict[str, Any]) -> tuple[str, dict[str, Any]]: + """Compute the recommended route and both conversion options for one component. + + Args: + members: The component's pipeline names. + inventory: The discover ``inventory.json`` document (deterministic or enriched). + + Returns: + ``(recommended, options)`` where ``recommended`` is ``"deterministic"`` when the whole + component is engine-capable, else ``"agentic"``; and ``options`` carries the ``deterministic`` + and ``agentic`` peers as first-class entries. + """ + deterministic = _deterministic_option(members, inventory) + agentic = _agentic_option(members, inventory) + recommended = _DETERMINISTIC_STRATEGY if deterministic["capable"] else "agentic" + return recommended, {"deterministic": deterministic, "agentic": agentic} + + +def build_recommendation(inventory: dict[str, Any]) -> dict[str, Any]: + """Compute the full routing recommendation over every component in the inventory. + + Returns a dict with ``components`` (each carrying ``component_id``, sorted ``members``, + ``recommended`` and both ``options``), the ``findings`` from component computation, and a + ``default_plan`` that proposes ``decision == recommended`` for every component -- ready to hand + straight to :func:`record_plan` when the user accepts the recommendations wholesale, or to edit + per component for overrides. + """ + components, findings = build_components(inventory) + component_entries: list[dict[str, Any]] = [] + default_plan_components: list[dict[str, Any]] = [] + for index, members in enumerate(components, start=1): + component_id = f"component-{index}" + recommended, options = recommend_component(members, inventory) + component_entries.append( + {"component_id": component_id, "members": members, "recommended": recommended, "options": options} + ) + default_plan_components.append({"component_id": component_id, "members": members, "decision": recommended}) + return { + "components": component_entries, + "findings": findings, + "default_plan": {"components": default_plan_components}, + } + + +# --------------------------------------------------------------------------- # +# Validation. All violations are collected (never fail-fast) so the authoring +# agent can fix every problem in one pass. +# --------------------------------------------------------------------------- # + + +def _components_by_id(inventory: dict[str, Any]) -> dict[str, list[str]]: + """Computed components keyed by their library id (``component-``).""" + components, _ = build_components(inventory) + return {f"component-{index}": members for index, members in enumerate(components, start=1)} + + +def validate_plan(raw: Any, inventory: dict[str, Any]) -> list[str]: + """Validate an authored conversion plan against the inventory's connected components. + + Returns a list of human-readable violation strings; an empty list means the plan is valid. Never + raises on a malformed payload. The rules: + + * only the authored top-level key ``components`` is allowed (library-owned keys are rejected with + a hint), and each component entry may carry only the authored fields (``recommended`` / + ``options`` are library-computed and rejected on input); + * ``decision`` must be one of the known routes and ``rationale`` (when present) a non-empty string; + * every member must be a real inventory pipeline, and a component's ``members`` must exactly match + one computed connected component -- so a decision can never split a component or span two; + * the plan is a **bijection** over components: every component is decided exactly once (no + partial plan, no duplicate/conflicting decisions). + """ + if not isinstance(raw, dict): + return [f"conversion plan must be a JSON object, got {type(raw).__name__}"] + + violations: list[str] = [] + for key in sorted(set(raw) - _PLAN_TOP_KEYS): + hint = " (set by the library, not the author)" if key in _LIBRARY_TOP_KEYS else "" + violations.append(f"unknown top-level key: {key!r}{hint}") + + components = raw.get("components") + if not isinstance(components, list): + violations.append("'components' must be a list") + return violations + + names = _pipeline_names(inventory) + computed_by_id = _components_by_id(inventory) + computed_by_members = {frozenset(members): component_id for component_id, members in computed_by_id.items()} + + decided_ids: list[str] = [] + for index, component in enumerate(components): + loc = f"components[{index}]" + matched_id = _validate_component_entry(component, loc, names, computed_by_members, violations) + if matched_id is not None: + decided_ids.append(matched_id) + + for component_id, count in Counter(decided_ids).items(): + if count > 1: + violations.append( + f"component {component_id!r} is decided {count} times; each component needs exactly one decision" + ) + decided = set(decided_ids) + for component_id, members in computed_by_id.items(): + if component_id not in decided: + violations.append(f"component {component_id!r} ({members}) has no decision; every component must be routed") + return violations + + +def _validate_component_entry( + component: Any, + loc: str, + names: set[str], + computed_by_members: dict[frozenset[str], str], + violations: list[str], +) -> str | None: + """Validate one authored component entry, appending problems; return the matched component id. + + Returns the computed ``component-`` id this entry decides when its ``members`` exactly match a + connected component (so the caller can enforce the bijection), else ``None``. + """ + if not isinstance(component, dict): + violations.append(f"{loc} must be an object") + return None + + for key in sorted(set(component) - _COMPONENT_AUTHORED_KEYS): + hint = " (set by the library, not the author)" if key in _COMPONENT_LIBRARY_KEYS else "" + violations.append(f"{loc}: unknown field {key!r}{hint}") + + decision = component.get("decision") + if decision not in DECISIONS: + allowed = ", ".join(repr(value) for value in DECISIONS) + violations.append(f"{loc}: 'decision' must be one of {{{allowed}}}, got {decision!r}") + + rationale = component.get("rationale") + if rationale is not None and (not isinstance(rationale, str) or not rationale.strip()): + violations.append(f"{loc}: 'rationale' must be a non-empty string when present") + + component_id = component.get("component_id") + if not isinstance(component_id, str) or not component_id: + violations.append(f"{loc}: 'component_id' must be a non-empty string") + + members = component.get("members") + if not isinstance(members, list) or not all(isinstance(member, str) for member in members): + violations.append(f"{loc}: 'members' must be a list of pipeline names") + return None + + for member in members: + if member not in names: + violations.append(f"{loc}: pipeline {member!r} not in inventory") + + # Reject duplicate members explicitly: a frozenset match would collapse ["a", "a"] to {"a"} and + # wrongly accept it as the component {"a"}, breaking the members-match / bijection contract. + duplicates = sorted({member for member in members if members.count(member) > 1}) + if duplicates: + violations.append(f"{loc}: duplicate members {duplicates}; list each pipeline once") + return None + + matched_id = computed_by_members.get(frozenset(members)) + if matched_id is None: + violations.append( + f"{loc}: members {sorted(members)} do not form a connected component " + f"(they split or span computed components); route each component as a whole" + ) + return None + if isinstance(component_id, str) and component_id and component_id != matched_id: + violations.append( + f"{loc}: component_id {component_id!r} does not match the component for these members " + f"(expected {matched_id!r})" + ) + return matched_id + + +# --------------------------------------------------------------------------- # +# Loading, recording, and the atomic idempotent write. +# --------------------------------------------------------------------------- # + + +def load_plan(*, plan: dict[str, Any] | None = None, plan_path: Path | None = None) -> dict[str, Any]: + """Return the raw authored plan dict from exactly one source (inline or file). + + Raises: + ValueError: if neither or both sources are provided. + """ + if (plan is None) == (plan_path is None): + raise ValueError("provide exactly one of 'plan' (inline dict) or 'plan_path'") + if plan is not None: + return plan + assert plan_path is not None # guaranteed by the guard above + return json.loads(Path(plan_path).read_text(encoding="utf-8")) + + +def build_plan_document(inventory: dict[str, Any], raw: dict[str, Any]) -> dict[str, Any]: + """Build the recorded plan document from a validated authored plan. + + The library recomputes ``members`` / ``recommended`` / both ``options`` and overlays only the + authored ``decision`` (and optional ``rationale``) per component, so the recorded facts cannot + drift from the inventory. Stamps the library-owned ``schema_version`` and ``inventory_sha256``. + Does not mutate the inputs and performs no I/O. + """ + recommendation = build_recommendation(inventory) + computed_by_id = _components_by_id(inventory) + members_to_id = {frozenset(members): component_id for component_id, members in computed_by_id.items()} + + authored_by_id: dict[str, dict[str, Any]] = {} + for component in raw.get("components", []): + component_id = members_to_id[frozenset(component["members"])] + authored_by_id[component_id] = component + + components_out: list[dict[str, Any]] = [] + for entry in recommendation["components"]: + component_id = entry["component_id"] + authored = authored_by_id[component_id] + recorded: dict[str, Any] = { + "component_id": component_id, + "members": entry["members"], + "recommended": entry["recommended"], + "decision": authored["decision"], + "options": entry["options"], + } + rationale = authored.get("rationale") + if rationale is not None: + recorded["rationale"] = rationale + components_out.append(recorded) + + return { + "schema_version": SCHEMA_VERSION, + "inventory_sha256": inventory_fingerprint(inventory), + "components": components_out, + "findings": recommendation["findings"], + } + + +def _write_plan_atomic(path: Path, document: dict[str, Any]) -> None: + """Write the plan JSON atomically (temp file + ``os.replace``), matching the inventory formatting.""" + path.parent.mkdir(parents=True, exist_ok=True) + temporary = path.with_name(f".{path.name}.tmp") + temporary.write_text(json.dumps(document, indent=2), encoding="utf-8") + os.replace(temporary, path) + + +def record_plan( + output_dir: Path, + *, + plan: dict[str, Any] | None = None, + plan_path: Path | None = None, +) -> dict[str, Any]: + """Validate an authored plan against the inventory, then record it on success. + + Reads ``/metadata/inventory.json``, validates the authored plan, and -- only when + there are no violations -- writes ``/metadata/conversion_plan.json`` atomically. The + inventory file is never touched. Provide the authored plan via exactly one of ``plan`` (inline + dict) or ``plan_path`` (a JSON file). + + Returns ``{"ok", "violations", "inventory_sha256", "components", "findings"}``. ``ok`` is + ``False`` (and no plan written) when there are violations. + + Raises: + FileNotFoundError: when ``inventory.json`` does not exist (run discover first). + ValueError: when neither or both plan sources are provided, or the inventory is not a JSON + object. + """ + inventory_path = Path(output_dir) / "metadata" / "inventory.json" + if not inventory_path.exists(): + raise FileNotFoundError(f"No inventory.json under {inventory_path.parent}; run the discover phase first.") + inventory = json.loads(inventory_path.read_text(encoding="utf-8")) + if not isinstance(inventory, dict): + raise ValueError(f"inventory.json must contain a JSON object, got {type(inventory).__name__}") + + raw = load_plan(plan=plan, plan_path=plan_path) + violations = validate_plan(raw, inventory) + if violations: + return {"ok": False, "violations": violations, "components": 0} + + document = build_plan_document(inventory, raw) + _write_plan_atomic(inventory_path.with_name(PLAN_FILENAME), document) + return { + "ok": True, + "violations": [], + "inventory_sha256": document["inventory_sha256"], + "components": len(document["components"]), + "findings": len(document["findings"]), + } diff --git a/tests/unit/test_mcp_route.py b/tests/unit/test_mcp_route.py new file mode 100644 index 0000000..b36e941 --- /dev/null +++ b/tests/unit/test_mcp_route.py @@ -0,0 +1,88 @@ +"""Tests that the MCP ``route`` command forwards to the adapter CLI correctly (routing, #77).""" + +from __future__ import annotations + +import json +from pathlib import Path +from typing import Any + +import pytest + +pytest.importorskip("mcp") + +from flowx.mcp import runner, server # noqa: E402 + + +class _Result: + ok = True + returncode = 0 + + def __init__(self, stdout: str) -> None: + self.stdout = stdout + self.stderr = "" + + def as_dict(self) -> dict[str, Any]: + return {"returncode": 0, "stdout": self.stdout, "stderr": ""} + + +@pytest.fixture +def captured(monkeypatch): + calls: list[list[str]] = [] + + def fake_run_adapter(args, **_kwargs): + calls.append([str(a) for a in args]) + if args and str(args[1]) == "recommend": + return _Result(json.dumps({"components": [{"component_id": "component-1"}], "default_plan": {}})) + return _Result(json.dumps({"ok": True, "violations": [], "components": 1})) + + monkeypatch.setattr(runner, "run_adapter", fake_run_adapter) + return calls + + +def test_route_recommend_forwards_to_adapter(captured, tmp_path: Path) -> None: + out = server._cmd_route({"output_dir": str(tmp_path), "action": "recommend"}) + assert out["ok"] is True + assert out["result"]["components"][0]["component_id"] == "component-1" + argv = captured[0] + assert argv[0] == "route" and argv[1] == "recommend" + assert "--output-dir" in argv + + +def test_route_defaults_to_recommend(captured, tmp_path: Path) -> None: + server._cmd_route({"output_dir": str(tmp_path)}) + assert captured[0][1] == "recommend" + + +def test_route_record_inline_plan_is_staged_to_a_file(captured, tmp_path: Path) -> None: + plan = {"components": [{"component_id": "component-1", "members": ["a"], "decision": "deterministic"}]} + out = server._cmd_route({"output_dir": str(tmp_path), "action": "record", "plan": plan}) + assert out["ok"] is True + argv = captured[0] + assert argv[0] == "route" and argv[1] == "record" + # The handler forwarded a real --plan-path (the staged temp file) to the CLI. + assert "--plan-path" in argv + + +def test_route_record_forwards_plan_path(captured, tmp_path: Path) -> None: + plan_path = tmp_path / "plan.json" + plan_path.write_text("{}", encoding="utf-8") + server._cmd_route({"output_dir": str(tmp_path), "action": "record", "plan_path": str(plan_path)}) + argv = captured[0] + assert argv[1] == "record" + assert argv[argv.index("--plan-path") + 1] == str(plan_path) + + +def test_route_record_requires_a_plan(tmp_path: Path) -> None: + both = server._cmd_route({"output_dir": str(tmp_path), "action": "record", "plan": {}, "plan_path": "x.json"}) + neither = server._cmd_route({"output_dir": str(tmp_path), "action": "record"}) + assert both["ok"] is False and "exactly one" in both["error"] + assert neither["ok"] is False and "exactly one" in neither["error"] + + +def test_route_rejects_unknown_action(tmp_path: Path) -> None: + out = server._cmd_route({"output_dir": str(tmp_path), "action": "sideways"}) + assert out["ok"] is False and "action" in out["error"] + + +def test_route_registered_in_command_map() -> None: + assert "route" in server._COMMANDS diff --git a/tests/unit/test_routing_components.py b/tests/unit/test_routing_components.py new file mode 100644 index 0000000..d5b458e --- /dev/null +++ b/tests/unit/test_routing_components.py @@ -0,0 +1,152 @@ +"""Tests for connected-component computation over control lineage (routing, #77). + +The inventory fixtures are built by hand through the source-agnostic emitter +(:func:`flowx.discovery_inventory.build_source_inventory`) -- no ADF, no Airflow -- so the +routing engine is proven against the standardised inventory shape and its per-pipeline +control-edge lineage, exactly as it will see it in production. Components are weak/undirected: +two pipelines land in the same component when a resolved control edge joins them in either +direction. +""" + +from __future__ import annotations + +from typing import Any + +from flowx.discovery_inventory import STRATEGY_PROPERTY, build_source_inventory +from flowx.models.discovery import CONCEPT_NOTEBOOK, SourceGraph, SourceNode +from flowx.models.ir import ControlEdge, Lineage +from flowx.routing import build_components + + +def _node(task_key: str, native_type: str = "Notebook", *, strategy: str = "deterministic") -> SourceNode: + return SourceNode( + source_id=task_key, + task_key=task_key, + concept=CONCEPT_NOTEBOOK, + source="unit", + name=task_key, + native_type=native_type, + properties={STRATEGY_PROPERTY: strategy}, + raw={"name": task_key, "type": native_type}, + ) + + +def _graph(name: str, *, edges: list[ControlEdge] | None = None, nodes: list[SourceNode] | None = None) -> SourceGraph: + """A one-node graph named *name*, optionally carrying control edges to other graphs.""" + return SourceGraph( + name=name, + source="unit", + tasks=nodes if nodes is not None else [_node(f"{name}_task")], + lineage=Lineage(control_edges=edges or []), + ) + + +def _edge(source: str, target: str, via: str, *, resolved: bool = True) -> ControlEdge: + return ControlEdge(source_workflow=source, target_workflow=target, via_task_key=via, resolved=resolved) + + +def _inventory(graphs: list[SourceGraph]) -> dict[str, Any]: + return build_source_inventory(graphs, source="unit", source_dir="/tmp/src") + + +def test_isolated_pipelines_each_form_their_own_component() -> None: + inventory = _inventory([_graph("a"), _graph("b"), _graph("c")]) + components, findings = build_components(inventory) + assert components == [["a"], ["b"], ["c"]] + assert findings == [] + + +def test_a_chain_of_calls_forms_one_component() -> None: + graphs = [ + _graph("a", edges=[_edge("a", "b", "call_b")], nodes=[_node("call_b", "ExecutePipeline")]), + _graph("b", edges=[_edge("b", "c", "call_c")], nodes=[_node("call_c", "ExecutePipeline")]), + _graph("c"), + ] + components, findings = build_components(_inventory(graphs)) + assert components == [["a", "b", "c"]] + assert findings == [] + + +def test_a_branch_fans_into_one_component() -> None: + graphs = [ + _graph( + "a", + edges=[_edge("a", "b", "call_b"), _edge("a", "c", "call_c")], + nodes=[_node("call_b", "ExecutePipeline"), _node("call_c", "ExecutePipeline")], + ), + _graph("b"), + _graph("c"), + ] + components, _ = build_components(_inventory(graphs)) + assert components == [["a", "b", "c"]] + + +def test_a_cycle_forms_one_component() -> None: + graphs = [ + _graph("a", edges=[_edge("a", "b", "call_b")], nodes=[_node("call_b", "ExecutePipeline")]), + _graph("b", edges=[_edge("b", "a", "call_a")], nodes=[_node("call_a", "ExecutePipeline")]), + ] + components, findings = build_components(_inventory(graphs)) + assert components == [["a", "b"]] + assert findings == [] + + +def test_multiple_edges_between_the_same_pair_still_one_component() -> None: + graphs = [ + _graph( + "a", + edges=[_edge("a", "b", "call_b1"), _edge("a", "b", "call_b2")], + nodes=[_node("call_b1", "ExecutePipeline"), _node("call_b2", "ExecutePipeline")], + ), + _graph("b"), + ] + components, _ = build_components(_inventory(graphs)) + assert components == [["a", "b"]] + + +def test_two_separate_clusters_are_distinct_components() -> None: + graphs = [ + _graph("a", edges=[_edge("a", "b", "call_b")], nodes=[_node("call_b", "ExecutePipeline")]), + _graph("b"), + _graph("x", edges=[_edge("x", "y", "call_y")], nodes=[_node("call_y", "ExecutePipeline")]), + _graph("y"), + ] + components, _ = build_components(_inventory(graphs)) + assert components == [["a", "b"], ["x", "y"]] + + +def test_unresolved_callee_is_a_finding_not_a_severed_edge() -> None: + graphs = [ + _graph( + "a", edges=[_edge("a", "", "call_ghost", resolved=False)], nodes=[_node("call_ghost", "ExecutePipeline")] + ) + ] + components, findings = build_components(_inventory(graphs)) + # 'a' still forms its own component; the unresolved edge is recorded, not dropped. + assert components == [["a"]] + assert len(findings) == 1 + assert "call_ghost" in findings[0] and "unresolved" in findings[0].lower() + + +def test_edge_to_pipeline_absent_from_inventory_is_a_finding() -> None: + graphs = [ + _graph("a", edges=[_edge("a", "ghost", "call_ghost")], nodes=[_node("call_ghost", "ExecutePipeline")]), + _graph("b"), + ] + components, findings = build_components(_inventory(graphs)) + assert components == [["a"], ["b"]] + assert len(findings) == 1 + assert "ghost" in findings[0] + + +def test_component_membership_and_ordering_are_deterministic() -> None: + # Declared out of alphabetical order; members and components must still sort deterministically. + graphs = [ + _graph("z"), + _graph("m", edges=[_edge("m", "a", "call_a")], nodes=[_node("call_a", "ExecutePipeline")]), + _graph("a"), + ] + first, _ = build_components(_inventory(graphs)) + second, _ = build_components(_inventory(graphs)) + assert first == second + assert first == [["a", "m"], ["z"]] diff --git a/tests/unit/test_routing_recommend.py b/tests/unit/test_routing_recommend.py new file mode 100644 index 0000000..43369b7 --- /dev/null +++ b/tests/unit/test_routing_recommend.py @@ -0,0 +1,206 @@ +"""Tests for the per-component conversion recommendation (routing, #77). + +Each component surfaces BOTH options as first-class peers: a deterministic option (engine-capability +assessment + motif/coverage evidence + uncovered gaps) and an agentic option (recommended Databricks +patterns from the insights block, with any simplification pattern surfaced prominently). The +``recommended`` field is a library-computed starting suggestion -- deterministic only when the whole +component is engine-capable. +""" + +from __future__ import annotations + +from typing import Any + +from flowx.discovery_inventory import STRATEGY_PROPERTY, build_source_inventory +from flowx.models.discovery import CONCEPT_NOTEBOOK, SourceGraph, SourceNode +from flowx.models.ir import ControlEdge, Lineage +from flowx.models.motifs import MOTIF_METADATA_DRIVEN_BULK_COPY, DetectedMotif +from flowx.routing import build_recommendation, recommend_component + + +def _node(task_key: str, native_type: str, *, strategy: str = "deterministic") -> SourceNode: + return SourceNode( + source_id=task_key, + task_key=task_key, + concept=CONCEPT_NOTEBOOK, + source="unit", + name=task_key, + native_type=native_type, + properties={STRATEGY_PROPERTY: strategy}, + raw={"name": task_key, "type": native_type}, + ) + + +def _inventory( + graphs: list[SourceGraph], + *, + motifs: dict[str, list[DetectedMotif]] | None = None, + insights: dict[str, Any] | None = None, +) -> dict[str, Any]: + inventory = build_source_inventory(graphs, source="unit", source_dir="/tmp/src", motifs_by_pipeline=motifs) + if insights is not None: + inventory["insights"] = insights + return inventory + + +def test_all_deterministic_component_recommends_deterministic() -> None: + graphs = [SourceGraph(name="a", source="unit", tasks=[_node("load", "Notebook"), _node("copy", "Copy")])] + recommended, options = recommend_component(["a"], _inventory(graphs)) + assert recommended == "deterministic" + assert options["deterministic"]["capable"] is True + assert options["deterministic"]["uncovered"] == [] + assert options["deterministic"]["activity_counts"] == {"deterministic": 2, "agentic": 0, "unsupported": 0} + + +def test_uncovered_agentic_activity_recommends_agentic_and_cites_the_gap() -> None: + graphs = [ + SourceGraph( + name="a", + source="unit", + tasks=[_node("load", "Notebook"), _node("flow", "ExecuteDataFlow", strategy="agentic")], + ) + ] + recommended, options = recommend_component(["a"], _inventory(graphs)) + assert recommended == "agentic" + assert options["deterministic"]["capable"] is False + assert options["deterministic"]["uncovered"] == [ + {"pipeline": "a", "activity": "flow", "type": "ExecuteDataFlow", "strategy": "agentic"} + ] + assert options["deterministic"]["activity_counts"] == {"deterministic": 1, "agentic": 1, "unsupported": 0} + + +def test_motif_covered_agentic_activity_is_engine_capable() -> None: + # An agentic-strategy activity that a detected motif claims is covered -> the component stays + # engine-capable and the recommendation is deterministic. + graphs = [ + SourceGraph( + name="a", + source="unit", + tasks=[_node("lookup", "Lookup"), _node("each", "ForEach", strategy="agentic")], + ) + ] + motif = DetectedMotif( + definition=MOTIF_METADATA_DRIVEN_BULK_COPY, + matched_activities=["lookup", "each"], + source_type_hint="database", + ) + recommended, options = recommend_component(["a"], _inventory(graphs, motifs={"a": [motif]})) + assert recommended == "deterministic" + assert options["deterministic"]["capable"] is True + assert options["deterministic"]["uncovered"] == [] + assert options["deterministic"]["motifs"] == ["metadata_driven_bulk_copy"] + + +def test_motif_coverage_does_not_leak_across_pipelines_in_a_component() -> None: + # A component spans A -> B. Pipeline A has a motif claiming an activity named 'shared'; pipeline B + # has an UNRELATED activity also named 'shared' with no motif. Motif coverage must be keyed by + # (pipeline, activity), so B's 'shared' stays an uncovered gap rather than being masked by A's. + graphs = [ + SourceGraph( + name="A", + source="unit", + tasks=[_node("call_b", "ExecutePipeline"), _node("shared", "Copy", strategy="agentic")], + lineage=Lineage( + control_edges=[ControlEdge(source_workflow="A", target_workflow="B", via_task_key="call_b")] + ), + ), + SourceGraph(name="B", source="unit", tasks=[_node("shared", "Copy", strategy="agentic")]), + ] + motif = DetectedMotif( + definition=MOTIF_METADATA_DRIVEN_BULK_COPY, matched_activities=["shared"], source_type_hint="database" + ) + recommended, options = recommend_component(["A", "B"], _inventory(graphs, motifs={"A": [motif]})) + assert recommended == "agentic" + assert options["deterministic"]["capable"] is False + # A's 'shared' is motif-covered; only B's unrelated 'shared' is the uncovered gap. + assert options["deterministic"]["uncovered"] == [ + {"pipeline": "B", "activity": "shared", "type": "Copy", "strategy": "agentic"} + ] + + +def test_missing_strategy_counts_as_unsupported_gap() -> None: + node = _node("mystery", "Custom") + node.properties.pop(STRATEGY_PROPERTY) # no strategy recorded at all + graphs = [SourceGraph(name="a", source="unit", tasks=[node])] + recommended, options = recommend_component(["a"], _inventory(graphs)) + assert recommended == "agentic" + assert options["deterministic"]["activity_counts"] == {"deterministic": 0, "agentic": 0, "unsupported": 1} + assert options["deterministic"]["uncovered"][0]["activity"] == "mystery" + + +def test_agentic_option_surfaces_insight_patterns_with_simplification_prominent() -> None: + graphs = [SourceGraph(name="a", source="unit", tasks=[_node("load", "Notebook")])] + insights = { + "pipeline_insights": [ + { + "pipeline": "a", + "recommended_patterns": [ + { + "pattern": "Lakeflow Connect SQL Server connector", + "fit": "Replaces the bespoke extractor", + "simplification_pattern": True, + }, + {"pattern": "Parameterised Lakeflow Job", "fit": "Like-for-like", "simplification_pattern": False}, + ], + } + ] + } + _, options = recommend_component(["a"], _inventory(graphs, insights=insights)) + agentic = options["agentic"] + assert agentic["has_simplification"] is True + assert agentic["recommended_patterns"] == [ + { + "pipeline": "a", + "pattern": "Lakeflow Connect SQL Server connector", + "fit": "Replaces the bespoke extractor", + "simplification_pattern": True, + }, + { + "pipeline": "a", + "pattern": "Parameterised Lakeflow Job", + "fit": "Like-for-like", + "simplification_pattern": False, + }, + ] + + +def test_agentic_option_empty_when_no_insights() -> None: + graphs = [SourceGraph(name="a", source="unit", tasks=[_node("load", "Notebook")])] + _, options = recommend_component(["a"], _inventory(graphs)) + assert options["agentic"] == {"recommended_patterns": [], "has_simplification": False} + + +def test_build_recommendation_emits_both_options_and_a_default_plan() -> None: + graphs = [ + SourceGraph( + name="parent", + source="unit", + tasks=[_node("call_child", "ExecutePipeline")], + lineage=Lineage( + control_edges=[ + ControlEdge(source_workflow="parent", target_workflow="child", via_task_key="call_child") + ] + ), + ), + SourceGraph(name="child", source="unit", tasks=[_node("flow", "ExecuteDataFlow", strategy="agentic")]), + SourceGraph(name="solo", source="unit", tasks=[_node("load", "Notebook")]), + ] + recommendation = build_recommendation(_inventory(graphs)) + ids = [component["component_id"] for component in recommendation["components"]] + assert ids == ["component-1", "component-2"] + first = recommendation["components"][0] + assert first["members"] == ["child", "parent"] + assert set(first["options"]) == {"deterministic", "agentic"} + assert first["recommended"] == "agentic" # child's ExecuteDataFlow is an uncovered gap + assert recommendation["components"][1]["recommended"] == "deterministic" # solo notebook + # The default plan proposes decision == recommended for every component, ready to record as-is. + assert recommendation["default_plan"]["components"] == [ + {"component_id": "component-1", "members": ["child", "parent"], "decision": "agentic"}, + {"component_id": "component-2", "members": ["solo"], "decision": "deterministic"}, + ] + + +def test_recommendation_is_pure_and_idempotent() -> None: + graphs = [SourceGraph(name="a", source="unit", tasks=[_node("load", "Notebook")])] + inventory = _inventory(graphs) + assert build_recommendation(inventory) == build_recommendation(inventory) diff --git a/tests/unit/test_routing_record.py b/tests/unit/test_routing_record.py new file mode 100644 index 0000000..8c3a1c1 --- /dev/null +++ b/tests/unit/test_routing_record.py @@ -0,0 +1,364 @@ +"""Tests for validating an authored conversion plan and recording it (routing, #77). + +The agent authors ONLY the per-component decision (and an optional rationale); the library recomputes +members, recommended, and both options' evidence on record, and binds the plan to the inventory with +a SHA-256 fingerprint. The recorded ``conversion_plan.json`` is a separate additive artifact -- it +never mutates ``inventory.json``. Fixtures are built through the source-agnostic emitter so the engine +is proven against the production inventory shape. +""" + +from __future__ import annotations + +import copy +import json +from pathlib import Path +from typing import Any + +import pytest + +from flowx.adapter.__main__ import main as adapter_cli_main +from flowx.discovery_insights import inventory_fingerprint +from flowx.discovery_inventory import STRATEGY_PROPERTY, build_source_inventory +from flowx.models.conversion_plan import ( + SCHEMA_VERSION, + ComponentPlan, + ConversionPlan, +) +from flowx.models.discovery import CONCEPT_NOTEBOOK, CONCEPT_RUN_WORKFLOW, SourceGraph, SourceNode +from flowx.models.ir import ControlEdge, Lineage +from flowx.routing import ( + PLAN_FILENAME, + build_recommendation, + load_plan, + record_plan, + validate_plan, +) + + +def _node(task_key: str, native_type: str, *, strategy: str = "deterministic") -> SourceNode: + return SourceNode( + source_id=task_key, + task_key=task_key, + concept=CONCEPT_NOTEBOOK, + source="unit", + name=task_key, + native_type=native_type, + properties={STRATEGY_PROPERTY: strategy}, + raw={"name": task_key, "type": native_type}, + ) + + +def _inventory() -> dict[str, Any]: + """A two-pipeline component (parent -> child, both engine-capable) plus a standalone 'solo'. + + 'child' carries a Lakeflow Connect simplification insight, so its component is deterministic-capable + AND has a prominent agentic re-architecture option -- the deterministic-vs-LFC choice #77 surfaces. + """ + parent = SourceGraph( + name="parent", + source="unit", + tasks=[_node("call_child", "ExecutePipeline")], + lineage=Lineage( + control_edges=[ControlEdge(source_workflow="parent", target_workflow="child", via_task_key="call_child")] + ), + ) + child = SourceGraph(name="child", source="unit", tasks=[_node("copy_orders", "Copy")]) + solo = SourceGraph(name="solo", source="unit", tasks=[_node("load", "Notebook")]) + inventory = build_source_inventory([parent, child, solo], source="unit", source_dir="/tmp/src") + inventory["insights"] = { + "pipeline_insights": [ + { + "pipeline": "child", + "recommended_patterns": [ + { + "pattern": "Lakeflow Connect SQL Server connector", + "fit": "Replaces the bespoke Copy extractor with a managed pipeline", + "simplification_pattern": True, + } + ], + } + ] + } + return inventory + + +def _authored_plan() -> dict[str, Any]: + """Accept-the-recommendation plan: decision == recommended for both components.""" + return { + "components": [ + {"component_id": "component-1", "members": ["child", "parent"], "decision": "deterministic"}, + {"component_id": "component-2", "members": ["solo"], "decision": "deterministic"}, + ] + } + + +def _write_inventory(output_dir: Path, inventory: dict[str, Any]) -> Path: + metadata = output_dir / "metadata" + metadata.mkdir(parents=True, exist_ok=True) + path = metadata / "inventory.json" + path.write_text(json.dumps(inventory, indent=2), encoding="utf-8") + return path + + +# --------------------------------------------------------------------------- # +# Models. +# --------------------------------------------------------------------------- # + + +def test_models_construct_and_document_the_contract() -> None: + plan = ConversionPlan( + inventory_sha256="abc", + components=[ + ComponentPlan( + component_id="component-1", + members=["child", "parent"], + recommended="deterministic", + decision="agentic", + ) + ], + ) + assert plan.schema_version == SCHEMA_VERSION + assert plan.components[0].decision == "agentic" + + +# --------------------------------------------------------------------------- # +# Validation: success (including the default plan round-trip). +# --------------------------------------------------------------------------- # + + +def test_valid_plan_passes_validation() -> None: + assert validate_plan(_authored_plan(), _inventory()) == [] + + +def test_default_plan_from_recommendation_validates() -> None: + inventory = _inventory() + default_plan = build_recommendation(inventory)["default_plan"] + assert validate_plan(default_plan, inventory) == [] + + +# --------------------------------------------------------------------------- # +# Validation: failure modes. +# --------------------------------------------------------------------------- # + + +def test_non_dict_payload_is_a_violation() -> None: + assert validate_plan([1, 2], _inventory()) == ["conversion plan must be a JSON object, got list"] + + +def test_unknown_top_level_key_including_library_owned_fields() -> None: + raw = _authored_plan() + raw["schema_version"] = "1" + raw["bogus"] = True + violations = validate_plan(raw, _inventory()) + assert any("unknown top-level key: 'schema_version' (set by the library, not the author)" in v for v in violations) + assert any("unknown top-level key: 'bogus'" in v for v in violations) + + +def test_library_owned_component_fields_are_rejected() -> None: + raw = _authored_plan() + raw["components"][0]["recommended"] = "deterministic" + raw["components"][0]["options"] = {} + violations = validate_plan(raw, _inventory()) + assert any("unknown field 'recommended' (set by the library, not the author)" in v for v in violations) + assert any("unknown field 'options' (set by the library, not the author)" in v for v in violations) + + +def test_invalid_decision_value_is_rejected() -> None: + raw = _authored_plan() + raw["components"][0]["decision"] = "maybe" + violations = validate_plan(raw, _inventory()) + assert any("'decision' must be one of" in v for v in violations) + + +def test_unknown_pipeline_in_members_is_rejected() -> None: + raw = _authored_plan() + raw["components"][0]["members"] = ["child", "ghost"] + violations = validate_plan(raw, _inventory()) + assert any("pipeline 'ghost' not in inventory" in v for v in violations) + + +def test_members_that_do_not_form_a_component_are_rejected() -> None: + # 'parent' and 'solo' are in different components; pairing them is not a coherent route. + raw = { + "components": [ + {"component_id": "component-1", "members": ["parent", "solo"], "decision": "deterministic"}, + ] + } + violations = validate_plan(raw, _inventory()) + assert any("do not form a connected component" in v for v in violations) + + +def test_partial_plan_missing_a_component_is_rejected() -> None: + raw = { + "components": [ + {"component_id": "component-1", "members": ["child", "parent"], "decision": "deterministic"}, + ] + } + violations = validate_plan(raw, _inventory()) + assert any("every component must be routed" in v for v in violations) + + +def test_duplicate_decision_for_one_component_is_rejected() -> None: + raw = { + "components": [ + {"component_id": "component-1", "members": ["child", "parent"], "decision": "deterministic"}, + {"component_id": "component-1", "members": ["child", "parent"], "decision": "agentic"}, + {"component_id": "component-2", "members": ["solo"], "decision": "deterministic"}, + ] + } + violations = validate_plan(raw, _inventory()) + assert any("decided" in v and "times" in v for v in violations) + + +def test_component_id_mismatch_is_rejected() -> None: + raw = _authored_plan() + raw["components"][0]["component_id"] = "component-9" + violations = validate_plan(raw, _inventory()) + assert any("does not match" in v for v in violations) + + +def test_duplicate_members_are_rejected() -> None: + # A frozenset would collapse ["solo", "solo"] to {"solo"} and wrongly match component-2; the + # members-match/bijection contract must reject the duplicate explicitly. + raw = _authored_plan() + raw["components"][1]["members"] = ["solo", "solo"] + violations = validate_plan(raw, _inventory()) + assert any("duplicate" in v.lower() and "solo" in v for v in violations) + + +def test_rationale_must_be_non_empty_when_present() -> None: + raw = _authored_plan() + raw["components"][0]["rationale"] = " " + violations = validate_plan(raw, _inventory()) + assert any("'rationale' must be a non-empty string" in v for v in violations) + + +# --------------------------------------------------------------------------- # +# Record: fingerprint binding, both-options shape, atomicity, idempotency. +# --------------------------------------------------------------------------- # + + +def test_record_writes_plan_with_both_options_recommended_and_fingerprint(tmp_path: Path) -> None: + inventory = _inventory() + _write_inventory(tmp_path, inventory) + result = record_plan(tmp_path, plan=_authored_plan()) + assert result["ok"] is True + assert result["components"] == 2 + + plan = json.loads((tmp_path / "metadata" / PLAN_FILENAME).read_text(encoding="utf-8")) + assert plan["schema_version"] == SCHEMA_VERSION + assert plan["inventory_sha256"] == inventory_fingerprint(inventory) + component = plan["components"][0] + assert component["members"] == ["child", "parent"] + assert component["recommended"] == "deterministic" + assert component["decision"] == "deterministic" + assert set(component["options"]) == {"deterministic", "agentic"} + # The agentic option's simplification pattern is surfaced as a first-class peer, not buried. + assert component["options"]["agentic"]["has_simplification"] is True + assert component["options"]["deterministic"]["capable"] is True + + +def test_record_preserves_a_user_override_with_both_recommended_and_decision(tmp_path: Path) -> None: + inventory = _inventory() + _write_inventory(tmp_path, inventory) + raw = _authored_plan() + # User overrides the deterministic recommendation to take the Lakeflow Connect re-architecture. + raw["components"][0]["decision"] = "agentic" + raw["components"][0]["rationale"] = "Adopt the managed connector to retire the extractor" + record_plan(tmp_path, plan=raw) + plan = json.loads((tmp_path / "metadata" / PLAN_FILENAME).read_text(encoding="utf-8")) + component = plan["components"][0] + assert component["recommended"] == "deterministic" + assert component["decision"] == "agentic" + assert component["rationale"] == "Adopt the managed connector to retire the extractor" + + +def test_record_leaves_inventory_byte_identical(tmp_path: Path) -> None: + inventory = _inventory() + inventory_path = _write_inventory(tmp_path, inventory) + original_bytes = inventory_path.read_bytes() + record_plan(tmp_path, plan=_authored_plan()) + assert inventory_path.read_bytes() == original_bytes + + +def test_record_is_idempotent(tmp_path: Path) -> None: + inventory = _inventory() + _write_inventory(tmp_path, inventory) + plan_path = tmp_path / "metadata" / PLAN_FILENAME + record_plan(tmp_path, plan=_authored_plan()) + first = plan_path.read_bytes() + record_plan(tmp_path, plan=_authored_plan()) + assert plan_path.read_bytes() == first + + +def test_record_leaves_plan_untouched_on_validation_failure(tmp_path: Path) -> None: + inventory = _inventory() + _write_inventory(tmp_path, inventory) + record_plan(tmp_path, plan=_authored_plan()) # write a good plan first + good_bytes = (tmp_path / "metadata" / PLAN_FILENAME).read_bytes() + + bad = _authored_plan() + bad["components"][0]["decision"] = "nonsense" + result = record_plan(tmp_path, plan=bad) + assert result["ok"] is False and result["violations"] + assert (tmp_path / "metadata" / PLAN_FILENAME).read_bytes() == good_bytes + + +def test_record_raises_when_inventory_missing(tmp_path: Path) -> None: + with pytest.raises(FileNotFoundError): + record_plan(tmp_path, plan=_authored_plan()) + + +def test_load_plan_requires_exactly_one_source(tmp_path: Path) -> None: + with pytest.raises(ValueError): + load_plan() + with pytest.raises(ValueError): + load_plan(plan={}, plan_path=tmp_path / "x.json") + + +# --------------------------------------------------------------------------- # +# CLI wiring. +# --------------------------------------------------------------------------- # + + +def test_cli_route_recommend_emits_components(tmp_path: Path, capsys: pytest.CaptureFixture[str]) -> None: + _write_inventory(tmp_path, _inventory()) + code = adapter_cli_main(["route", "recommend", "--output-dir", str(tmp_path)]) + assert code == 0 + payload = json.loads(capsys.readouterr().out) + assert [c["component_id"] for c in payload["components"]] == ["component-1", "component-2"] + assert "default_plan" in payload + + +def test_cli_route_record_success(tmp_path: Path, capsys: pytest.CaptureFixture[str]) -> None: + _write_inventory(tmp_path, _inventory()) + plan_path = tmp_path / "plan.json" + plan_path.write_text(json.dumps(_authored_plan()), encoding="utf-8") + code = adapter_cli_main(["route", "record", "--output-dir", str(tmp_path), "--plan-path", str(plan_path)]) + assert code == 0 + payload = json.loads(capsys.readouterr().out) + assert payload["ok"] is True and payload["components"] == 2 + assert (tmp_path / "metadata" / PLAN_FILENAME).exists() + + +def test_cli_route_record_validation_failure_returns_1(tmp_path: Path, capsys: pytest.CaptureFixture[str]) -> None: + _write_inventory(tmp_path, _inventory()) + raw = _authored_plan() + raw["components"].pop() # partial plan: component-2 undecided + plan_path = tmp_path / "plan.json" + plan_path.write_text(json.dumps(raw), encoding="utf-8") + code = adapter_cli_main(["route", "record", "--output-dir", str(tmp_path), "--plan-path", str(plan_path)]) + assert code == 1 + payload = json.loads(capsys.readouterr().out) + assert payload["ok"] is False and payload["violations"] + + +def test_fixture_isolation() -> None: + a = _authored_plan() + b = _authored_plan() + a["components"][0]["decision"] = "agentic" + assert b["components"][0]["decision"] == "deterministic" + assert copy.deepcopy(a) == a + # CONCEPT_RUN_WORKFLOW is imported for symmetry with the discovery fixtures; touch it so linters + # see the dependency used. + assert isinstance(CONCEPT_RUN_WORKFLOW, str)