diff --git a/skills/flowx-convert/sources/adf.md b/skills/flowx-convert/sources/adf.md index 2a605f33..471cec52 100644 --- a/skills/flowx-convert/sources/adf.md +++ b/skills/flowx-convert/sources/adf.md @@ -75,6 +75,58 @@ Write one JSON file per resolved gap into `/agentic_results/`: `activity_name` (required) is matched by name, recursing into containers. `task_key`/`depends_on` are inherited from the placeholder when omitted, preserving dependency edges. +When no typed task fits, use `"type": "AgenticComponentActivity"` instead. It carries `files` +(each a relative `path` below the bundle's `src/` with exactly one of text `content` or base64 +`binary_content`, written byte for byte), +pipeline `resources` (`resource_key` plus a `definition` mapping), optional job `environments` +(each exactly an `environment_key` plus a `spec` mapping), and a raw Databricks `task` fragment with +exactly one `_task` payload. A `spark_python_task`, `python_wheel_task`, `spark_jar_task` or +`dbt_task` must name its compute (`environment_key`, `job_cluster_key`, `existing_cluster_id` or +`new_cluster`), and a `spark_submit_task`, which runs only on a new cluster, must name `new_cluster` +or a `job_cluster_key`; otherwise package refuses it. A component cannot declare a job +cluster, so a `job_cluster_key` must be one flowx defines (`default_cluster`, `single_node_cluster` +or `multi_node_cluster`); flowx adds that cluster to the job. A serverless task uses an +`environment_key`: declare that environment in `environments` and flowx writes it into the job that +runs the task. A wheel or other authored file the environment installs is listed in +`spec.dependencies` as `../src/`. Every `environment_key` a task uses must be declared, and +two components may declare the same environment only with an identical spec (it is then written +once). + +flowx owns `task_key`, `depends_on`, `run_if`, `timeout_seconds`, `max_retries`, +`min_retry_interval_millis` and `retry_on_timeout`; setting any of them in the fragment is an +error. On merge, the component's `name`, `task_key`, `depends_on`, +`timeout_seconds`, `max_retries` and `min_retry_interval_millis` always come from the placeholder; +any values you set for them on the activity are replaced. flowx does not override or remove any +other value the fragment sets, except in three wiring passes that apply to authored tasks exactly as +to generated ones: + +- A `{{tasks.X.values.Y}}` reference to a task outside the same job is blanked, because task values + do not cross job boundaries. Only notebook `base_parameters`, `run_job_task.job_parameters` and + `condition_task` operands are checked. SETUP.md lists the blanked notebook parameters and + neutralised conditions, not blanked `job_parameters`; references in other payloads are left as + they are. +- Inside a ForEach that runs as its own job, notebook `base_parameters` and `condition_task` + operands are rewritten for that job: `{{input.x}}` becomes `{{job.parameters.x}}`, ADF + expressions become job-parameter references, and non-string values become strings. +- A `run_job_task` whose `job_id` is `${resources.jobs.X.id}` for a job outside this bundle is + pointed at an `X_job_id` bundle variable instead. + +Where the fragment leaves a value out, flowx adds only the plumbing the task needs: + +- The ForEach `item` parameter, only to a `notebook_task` (`base_parameters`) or `run_job_task` + (`job_parameters`) that is the ForEach's only child; any other payload of an only child passes + `{{input}}` itself. When the ForEach has several children, or its only child is an IfCondition or + Switch holding the component, the children run as a `_inner_tasks` job where `{{input}}` + does not resolve: reference `{{job.parameters.item}}` in a notebook `base_parameters`, a + `run_job_task.job_parameters` or a `condition_task` operand, which flowx forwards. Other payloads + are not scanned, so they cannot receive the item on their own. +- The source activity's collapsed notifications, when the fragment sets neither + `email_notifications` nor `webhook_notifications`. +- A job-cluster binding, only for a notebook task that names no compute and either has a + classic compute mode, ships libraries, or (without a serverless compute mode) points at a + workspace path outside the bundle. A bundle `../src/` notebook without libraries runs on + serverless, and no other task type is given a cluster, so name the compute you need. + ## Step 6 — Merge agentic results ```bash diff --git a/src/flowx/bundler/dab_writer.py b/src/flowx/bundler/dab_writer.py index c72fe5ad..f934943a 100644 --- a/src/flowx/bundler/dab_writer.py +++ b/src/flowx/bundler/dab_writer.py @@ -23,13 +23,14 @@ SINGLE_NODE_JOB_CLUSTER_KEY, ) from flowx.bundler.inner_job_params import normalize_value -from flowx.bundler.notebook_writer import write_notebooks +from flowx.bundler.notebook_writer import content_signature, write_notebooks from flowx.bundler.prereqs_writer import ManualParameter, build_prereqs, render_setup_md from flowx.bundler.setup_generator import generate_setup_tasks from flowx.ir_serde import data_asset_from_dict, lineage_from_dict from flowx.models.dab import DabNotebook from flowx.models.ir import ( Activity, + AgenticComponentActivity, AppendVariableActivity, CopyActivity, DbtFactoryActivity, @@ -112,6 +113,7 @@ def write_bundle( _cross_bundle_variables.clear() _neutralized_conditions.clear() + _check_agentic_file_collisions(workflow, catalog, schema) workflow = copy.deepcopy(workflow) output_dir = Path(output_dir) @@ -133,6 +135,9 @@ def write_bundle( ) pipeline_resources = _collect_pipeline_resources(workflow) + _check_pipeline_resource_keys(pipeline_resources, known_bundle_jobs) + environments = _collect_environments(workflow) + _check_environment_keys(environments, workflow) pipeline_variable_declarations = _build_pipeline_variable_declarations(pipeline_resources, catalog, schema) # sql_task references ${var.warehouse_id}; declare it (no default -> user supplies at deploy). if _bundle_uses_sql_task(workflow): @@ -183,7 +188,9 @@ def write_bundle( resources_dir.mkdir(parents=True, exist_ok=True) job_yml_path = resources_dir / f"{resource_key}.yml" hoisted_global_names = set(hoisted_global_variables) - job_resource = _build_job_resource(workflow, resource_key, hoisted_globals=hoisted_global_names) + job_resource = _build_job_resource( + workflow, resource_key, hoisted_globals=hoisted_global_names, environments=environments + ) job_yml_path.write_text( yaml.dump( job_resource, default_flow_style=False, sort_keys=False, allow_unicode=True, Dumper=_BundleYamlDumper @@ -198,7 +205,11 @@ def write_bundle( inner_key = normalize_task_key(inner.name) inner_yml_path = resources_dir / f"{inner_key}.yml" inner_resource = _build_job_resource( - inner, inner_key, extra_notebooks_for_augment=workflow.notebooks, hoisted_globals=hoisted_global_names + inner, + inner_key, + extra_notebooks_for_augment=workflow.notebooks, + hoisted_globals=hoisted_global_names, + environments=environments, ) inner_yml_path.write_text( yaml.dump( @@ -234,11 +245,7 @@ def write_bundle( src_dir = output_dir / "src" def _write_generated(notebooks: list[DabNotebook]) -> None: - root_artifacts = [ - notebook - for notebook in notebooks - if notebook.relative_path.startswith("resources/") or notebook.relative_path == "pyproject.toml" - ] + root_artifacts = [notebook for notebook in notebooks if _is_bundle_root_artifact(notebook)] rest = [notebook for notebook in notebooks if notebook not in root_artifacts] if rest: created_files.extend(write_notebooks(rest, src_dir)) @@ -656,6 +663,43 @@ def _known_bundle_job_keys(workflow: PreparedWorkflow, resource_key: str) -> set return keys +def _check_agentic_file_collisions(workflow: PreparedWorkflow, catalog: str, schema: str) -> None: + """Refuse an agent-authored file whose path another file below ``src`` uses with different content. + + The notebook writer gives a clashing file a ``__N`` suffix, and setup notebooks are written later and + replace whatever is there, but every task still points at the original path, so one task would run the + other file's code. + """ + authored_paths: set[str] = set() + signatures_by_path: dict[str, set[object]] = {} + for prepared in [workflow, *workflow.inner_workflows]: + setup_notebooks = generate_setup_tasks( + secrets=prepared.secrets, setup_tasks=prepared.setup_tasks, catalog=catalog, schema=schema + ) + for notebook in [*prepared.notebooks, *setup_notebooks]: + if notebook.authored: + authored_paths.add(notebook.relative_path) + elif _is_bundle_root_artifact(notebook): + continue + signatures_by_path.setdefault(notebook.relative_path, set()).add(content_signature(notebook)) + for relative_path in sorted(authored_paths): + if len(signatures_by_path[relative_path]) > 1: + raise ValueError( + f"Agentic component file {relative_path!r} shares its path with another file in the bundle that " + "has different content; give each component its own file path" + ) + + +def _is_bundle_root_artifact(notebook: DabNotebook) -> bool: + """Return whether a generated file belongs at the bundle root instead of below ``src``. + + Authored files always stay below ``src``, even when their path looks like a PyDABs artifact. + """ + if notebook.authored: + return False + return notebook.relative_path.startswith("resources/") or notebook.relative_path == "pyproject.toml" + + def _namespace_workflow_assets(workflow: PreparedWorkflow) -> PreparedWorkflow: """Namespaces generated source files by DAG while preserving workspace paths.""" cloned = copy.deepcopy(workflow) @@ -692,7 +736,7 @@ def _namespace_workflow_assets(workflow: PreparedWorkflow) -> PreparedWorkflow: notebook.content = notebook.content.replace("src/dbt_project", f"src/{prefix}/dbt_project") notebook.content = notebook.content.replace("src/dbt_profiles", f"src/{prefix}/dbt_profiles") continue - if original_path.startswith("resources/") or original_path == "pyproject.toml": + if _is_bundle_root_artifact(notebook): continue notebook.relative_path = f"{prefix}/{original_path}" replacements[f"../src/{original_path}"] = f"../src/{notebook.relative_path}" @@ -1049,14 +1093,97 @@ def _collect_pipeline_resources(workflow: PreparedWorkflow) -> list[dict[str, An Returns: Flat list of pipeline-resource dicts (each with ``resource_key`` - and ``definition``), including entries from inner workflows. + and ``definition``), including entries from inner workflows. A + resource declared more than once with the same key and definition + appears once, so it is written once; other fields on the entry are + never written and play no part in that comparison. """ - resources = list(workflow.pipeline_resources) - for inner in workflow.inner_workflows: - resources.extend(inner.pipeline_resources) + resources: list[dict[str, Any]] = [] + for current in [workflow, *workflow.inner_workflows]: + for resource in current.pipeline_resources: + if not any( + existing["resource_key"] == resource["resource_key"] + and existing["definition"] == resource["definition"] + for existing in resources + ): + resources.append(resource) return resources +def _check_pipeline_resource_keys(pipeline_resources: list[dict[str, Any]], job_resource_keys: set[str]) -> None: + """Fail when a pipeline resource key clashes with another resource in the bundle. + + Every pipeline resource is written to ``resources/.yml``, as is every static job, so a + pipeline key that matches a job key or an earlier pipeline key would silently replace that file + and drop the other resource. Bundle resource keys must also be unique across resource types, so + a match with a Python-generated dbt-factory job would fail at deploy time instead. Repeats with + the same definition are already collapsed by :func:`_collect_pipeline_resources`, so a repeated + key here always carries a different definition. + + Raises: + ValueError: A pipeline resource key matches a job resource key or repeats an earlier one. + """ + seen_pipeline_keys: set[str] = set() + for resource in pipeline_resources: + pipeline_key = resource["resource_key"] + if pipeline_key in job_resource_keys: + raise ValueError( + f"Pipeline resource key {pipeline_key!r} matches a job resource key in this bundle; " + "bundle resource keys must be unique" + ) + if pipeline_key in seen_pipeline_keys: + raise ValueError( + f"Pipeline resource key {pipeline_key!r} is used by more than one pipeline resource with " + f"different definitions; each would overwrite resources/{pipeline_key}.yml" + ) + seen_pipeline_keys.add(pipeline_key) + + +def _collect_environments(workflow: PreparedWorkflow) -> list[dict[str, Any]]: + """Returns every job environment carried by *workflow* and its inner jobs. + + Each entry is an ``{environment_key, spec}`` dict. An environment declared + more than once with the same key and spec appears once, so it is written + once per job that uses it. + """ + environments: list[dict[str, Any]] = [] + for current in [workflow, *workflow.inner_workflows]: + for environment in current.environments: + if environment not in environments: + environments.append(environment) + return environments + + +def _check_environment_keys(environments: list[dict[str, Any]], workflow: PreparedWorkflow) -> None: + """Fail when an environment key is declared with two specs or a task names an undeclared one. + + Each job lists the environments its tasks reference, so two specs under one key would leave a + task's compute ambiguous, and a task whose ``environment_key`` matches no declared environment is + rejected by the Jobs API at deploy time. Identical repeats are already collapsed by + :func:`_collect_environments`, so a repeated key here always carries a different spec. + + Raises: + ValueError: An environment key repeats with a different spec, or a task references an + environment key no environment declares. + """ + declared_keys: set[str] = set() + for environment in environments: + environment_key = environment["environment_key"] + if environment_key in declared_keys: + raise ValueError( + f"Job environment {environment_key!r} is declared more than once with different specs; " + "give each spec its own environment_key" + ) + declared_keys.add(environment_key) + for current in [workflow, *workflow.inner_workflows]: + for task in _iter_tasks_recursively(current.tasks): + if "environment_key" in task and task["environment_key"] not in declared_keys: + raise ValueError( + f"Task {task.get('task_key')!r} uses environment_key {task['environment_key']!r}, which no " + "agentic component declares in its environments" + ) + + def _collect_pydabs_resource_entries(workflow: PreparedWorkflow) -> list[str]: """Returns the ``python.resources`` entries for every dbt-factory PyDABs hook in *workflow*. @@ -1187,14 +1314,15 @@ def _collect_required_cluster_keys(tasks: list[dict[str, Any]]) -> set[str]: return {task["job_cluster_key"] for task in _iter_tasks_recursively(tasks) if task.get("job_cluster_key")} -def _strip_compute_mode_markers(tasks: list[dict[str, Any]]) -> None: - """Removes the private ``_compute_mode`` marker from every task before YAML output. +def _strip_private_task_markers(tasks: list[dict[str, Any]]) -> None: + """Removes the private ``_compute_mode`` and ``_authored`` markers from every task before YAML output. Args: tasks: Top-level task dicts (mutated in place). """ for task in _iter_tasks_recursively(tasks): task.pop("_compute_mode", None) + task.pop("_authored", None) # Patterns that signal a base_parameter value couldn't be evaluated cleanly. When a task references an @@ -1224,10 +1352,13 @@ def _extract_manual_parameters_from_existing_notebook_tasks( deterministic translation are emitted with raw ADF expression values that ``dbutils.widgets.get`` returns verbatim, which fails at runtime. Walking both the absolute-path and bundle-relative cases drops the - broken values and surfaces them as a SETUP.md row instead. + broken values and surfaces them as a SETUP.md row instead. Agent-authored + tasks (marked ``_authored``) are skipped: their parameters are written as given. """ manual_parameters: list[ManualParameter] = [] for task in _iter_tasks_recursively(tasks): + if task.get("_authored"): + continue notebook_task = task.get("notebook_task") or {} notebook_path = notebook_task.get("notebook_path", "") base_params = notebook_task.get("base_parameters") @@ -1263,7 +1394,7 @@ def _iter_tasks_recursively(tasks: list[dict[str, Any]]) -> Iterator[dict[str, A yield from _iter_tasks_recursively([inner]) -_CLUSTER_BINDING_KEYS = ("existing_cluster_id", "new_cluster", "job_cluster_key") +_COMPUTE_BINDING_KEYS = ("existing_cluster_id", "new_cluster", "job_cluster_key", "environment_key") def _any_task_uses_classic_cluster(tasks: list[dict[str, Any]]) -> bool: @@ -1286,6 +1417,8 @@ def _bind_cluster_to_notebook_tasks(tasks: list[dict[str, Any]]) -> None: matching job_cluster. Tasks without a marker fall back to the legacy behaviour: existing-workspace notebooks bind to ``default_cluster`` and flowx-generated notebooks stay unbound. + A task that already names its compute (a cluster or a serverless + ``environment_key``) is never given a second binding. Args: tasks: Top-level task dicts (mutated in place). @@ -1294,7 +1427,7 @@ def _bind_cluster_to_notebook_tasks(tasks: list[dict[str, Any]]) -> None: notebook_task = task.get("notebook_task") if notebook_task is None: continue - if any(key in task for key in _CLUSTER_BINDING_KEYS): + if any(key in task for key in _COMPUTE_BINDING_KEYS): continue compute_mode = task.get("_compute_mode") if compute_mode == "serverless": @@ -1563,14 +1696,16 @@ def _augment_base_parameters( Args: tasks: Top-level task dicts (mutated in place). - notebooks: Generated notebooks to scan. + notebooks: Generated notebooks to scan. Authored files are skipped so + their tasks keep exactly the parameters the agent wrote and the + notebook's own widget defaults still apply. hoisted_globals: Names of factory globals hoisted to bundle variables. A widget matching one of these binds to ``${var.NAME}`` so the deploy-time bundle variable flows into the notebook; other widgets default to an empty string as before. """ hoisted = hoisted_globals or set() - notebook_by_relpath = {notebook.relative_path: notebook for notebook in notebooks} + notebook_by_relpath = {notebook.relative_path: notebook for notebook in notebooks if not notebook.authored} def visit(task: dict[str, Any]) -> None: notebook_task = task.get("notebook_task") @@ -1599,6 +1734,7 @@ def _build_job_resource( attach_clusters: bool = True, extra_notebooks_for_augment: list[DabNotebook] | None = None, hoisted_globals: set[str] | None = None, + environments: list[dict[str, Any]] | None = None, ) -> dict[str, Any]: """Builds a job resource dict for a single workflow. @@ -1609,6 +1745,9 @@ def _build_job_resource( and binds every notebook task to it. Set to ``False`` for inner jobs that are invoked via ``run_job_task`` from another bundle job — they inherit compute from the caller. + environments: The bundle's job environments; the ones this job's + tasks reference by ``environment_key`` are written as its + ``environments`` block. Returns: Dict ready for YAML serialization. @@ -1650,7 +1789,18 @@ def _build_job_resource( extras=cluster_extras or None, ) - _strip_compute_mode_markers(workflow.tasks) + referenced_environment_keys = { + task["environment_key"] for task in _iter_tasks_recursively(workflow.tasks) if "environment_key" in task + } + job_environments = [ + environment + for environment in environments or [] + if environment["environment_key"] in referenced_environment_keys + ] + if job_environments: + job_def["environments"] = job_environments + + _strip_private_task_markers(workflow.tasks) if workflow.parameters: # Emit each job parameter once in the DAB shape ({name, default}); dropping the internal ``type`` @@ -2276,6 +2426,15 @@ def _reconstruct_ir(task_ir: dict[str, Any]) -> Activity: consolidate_metadata_driven=bool(task_ir.get("consolidate_metadata_driven", False)), lookup_values=list(task_ir.get("lookup_values") or []), ) + if task_type == "AgenticComponentActivity": + return AgenticComponentActivity( + **base, + files=list(task_ir.get("files") or []), + resources=list(task_ir.get("resources") or []), + environments=list(task_ir.get("environments") or []), + task=dict(task_ir.get("task") or {}), + raw_definition=task_ir.get("raw_definition"), + ) if task_type == "UnsupportedActivity": return UnsupportedActivity( **base, diff --git a/src/flowx/bundler/notebook_writer.py b/src/flowx/bundler/notebook_writer.py index 05285bb6..fc6ae901 100644 --- a/src/flowx/bundler/notebook_writer.py +++ b/src/flowx/bundler/notebook_writer.py @@ -10,7 +10,7 @@ logger = logging.getLogger(__name__) -def _content_signature(notebook: DabNotebook) -> object: +def content_signature(notebook: DabNotebook) -> object: """Return a value suitable for comparing two notebooks for byte-equivalence.""" if notebook.binary_content is not None: return ("binary", notebook.binary_content) @@ -50,7 +50,7 @@ def write_notebooks(notebooks: list[DabNotebook], output_dir: Path) -> list[Path for notebook in notebooks: target = notebook.relative_path - signature = _content_signature(notebook) + signature = content_signature(notebook) existing = written_signatures.get(target) if existing is not None: if existing == signature: @@ -70,6 +70,9 @@ def write_notebooks(notebooks: list[DabNotebook], output_dir: Path) -> list[Path destination.parent.mkdir(parents=True, exist_ok=True) if notebook.binary_content is not None: destination.write_bytes(notebook.binary_content) + elif notebook.authored: + # Authored text is written byte for byte: no added newline and no platform line endings. + destination.write_bytes(notebook.content.encode("utf-8")) else: content = notebook.content if notebook.content.endswith("\n") else notebook.content + "\n" destination.write_text(content, encoding="utf-8") diff --git a/src/flowx/ir_serde.py b/src/flowx/ir_serde.py index 3dbbe319..9b88419a 100644 --- a/src/flowx/ir_serde.py +++ b/src/flowx/ir_serde.py @@ -21,6 +21,7 @@ from flowx.models.ir import ( Activity, + AgenticComponentActivity, AppendVariableActivity, ControlEdge, CopyActivity, @@ -296,6 +297,14 @@ def activity_extra_fields(activity: Activity) -> dict[str, Any]: extra: dict[str, Any] = {} match activity: + case AgenticComponentActivity(): + extra["files"] = activity.files + extra["resources"] = activity.resources + if activity.environments: + extra["environments"] = activity.environments + extra["task"] = activity.task + if activity.raw_definition is not None: + extra["raw_definition"] = activity.raw_definition case NotebookActivity(): extra["notebook_path"] = activity.notebook_path if activity.base_parameters: @@ -576,21 +585,42 @@ def pipeline_to_debug_dict(pipeline: Pipeline) -> dict[str, Any]: } +_PLACEHOLDER_OWNED_COMPONENT_FIELDS = ( + "task_key", + "name", + "depends_on", + "timeout_seconds", + "max_retries", + "min_retry_interval_millis", +) + + def _find_and_replace_task(tasks: list[dict[str, Any]], activity_name: str, replacement: dict[str, Any]) -> bool: """Replace the task named *activity_name* with *replacement*, recursing into containers. Searches top-level tasks and the nested activity lists of IfCondition / ForEach / Switch containers. Preserves the placeholder's ``task_key`` and ``depends_on`` when the replacement omits them so downstream dependency - edges stay intact. Returns True when a match was replaced. + edges stay intact. An ``AgenticComponentActivity`` replacement instead + always takes its identity, dependencies, timeout, and retries from the + placeholder, dropping any value the agent set (or the placeholder lacks), + because flowx owns those for an authored component. Returns True when a + match was replaced. """ nested_keys = ("inner_activities", "if_true_activities", "if_false_activities", "default_activities") for index, task in enumerate(tasks): if task.get("name") == activity_name: - replacement.setdefault("task_key", task.get("task_key")) - replacement.setdefault("name", activity_name) - if "depends_on" not in replacement and task.get("depends_on"): - replacement["depends_on"] = task["depends_on"] + if replacement.get("type") == "AgenticComponentActivity": + for field_name in _PLACEHOLDER_OWNED_COMPONENT_FIELDS: + if field_name in task: + replacement[field_name] = task[field_name] + else: + replacement.pop(field_name, None) + else: + replacement.setdefault("task_key", task.get("task_key")) + replacement.setdefault("name", activity_name) + if "depends_on" not in replacement and task.get("depends_on"): + replacement["depends_on"] = task["depends_on"] tasks[index] = replacement return True for key in nested_keys: @@ -618,7 +648,9 @@ def merge_agentic_results(report_path: Path, results_dir: Path, output_path: Pat The matching placeholder task (located by ``name``, recursing into IfCondition / ForEach / Switch containers) is replaced by ``task``. Use a ``NotebookActivity`` whose ``notebook_path`` points at a notebook the agent - wrote to the workspace; the prepare phase then references it directly. + wrote to the workspace; the prepare phase then references it directly. An + ``AgenticComponentActivity`` keeps the placeholder's name, task key, + dependencies, timeout, and retries whatever the agent wrote. Args: report_path: ``translation_report.json`` produced by the translate phase. diff --git a/src/flowx/models/dab.py b/src/flowx/models/dab.py index ec0b85ce..cf659568 100644 --- a/src/flowx/models/dab.py +++ b/src/flowx/models/dab.py @@ -20,12 +20,16 @@ class DabNotebook: language: Notebook language (``"python"``, ``"sql"``, ``"scala"``, ``"r"``). binary_content: Raw bytes for binary files (e.g. JARs). When set, the notebook writer writes these bytes instead of ``content``. + authored: Marks a file an agent authored for an agentic component. The + bundle writer keeps it below ``src`` whatever its path and writes it + as given. """ relative_path: str content: str = "" language: str = "python" binary_content: bytes | None = None + authored: bool = False # --------------------------------------------------------------------------- diff --git a/src/flowx/models/ir.py b/src/flowx/models/ir.py index 6c98f635..10e683fc 100644 --- a/src/flowx/models/ir.py +++ b/src/flowx/models/ir.py @@ -688,6 +688,39 @@ class DbtFactoryActivity(Activity): nodes: list[dict[str, Any]] = field(default_factory=list) +@dataclass(slots=True, kw_only=True) +class AgenticComponentActivity(Activity): + """Bundle components authored for a source activity the typed engine cannot express. + + Attributes: + files: Files to write below the bundle's ``src`` directory. Each entry + carries a relative ``path`` and either UTF-8 ``content`` or + base64-encoded ``binary_content``. A file may share its path with + another component's file or a generated notebook only when the + content is identical. + resources: Pipeline resources in the existing ``resource_key`` plus + raw ``definition`` shape used by the bundle writer. Components may + declare the same resource only with an identical definition; it + is then written once. + environments: Job environments in the Jobs API ``environment_key`` + plus ``spec`` shape, for a serverless task that names one through + its ``environment_key``. The bundle writer adds each to every job + holding a task that references it. Components may declare the + same environment only with an identical spec; it is then written + once. + task: Raw Databricks task fragment carrying exactly one executable + payload (a ``_task`` key such as ``pipeline_task`` or + ``notebook_task``) wired to an authored resource or file. + raw_definition: Original source definition retained for auditing. + """ + + files: list[dict[str, Any]] = field(default_factory=list) + resources: list[dict[str, Any]] = field(default_factory=list) + environments: list[dict[str, Any]] = field(default_factory=list) + task: dict[str, Any] = field(default_factory=dict) + raw_definition: dict[str, Any] | None = None + + @dataclass(slots=True, kw_only=True) class UnsupportedActivity(Activity): """Sentinel for activities that could not be translated. diff --git a/src/flowx/preparer/activity_preparers/agentic_component.py b/src/flowx/preparer/activity_preparers/agentic_component.py new file mode 100644 index 00000000..cf91e22c --- /dev/null +++ b/src/flowx/preparer/activity_preparers/agentic_component.py @@ -0,0 +1,247 @@ +"""Lower agent-authored bundle components through the existing output channels.""" + +from __future__ import annotations + +import base64 +from pathlib import PurePosixPath, PureWindowsPath + +from flowx.bundler.constants import DEFAULT_JOB_CLUSTER_KEY, MULTI_NODE_JOB_CLUSTER_KEY, SINGLE_NODE_JOB_CLUSTER_KEY +from flowx.models.dab import DabNotebook +from flowx.models.ir import AgenticComponentActivity +from flowx.preparer.workflow_preparer import PreparedActivity, build_common_task_fields +from flowx.utils import normalize_task_key + +# flowx owns a task's identity, its upstream wiring, and its run policy; these come from +# the activity itself, so an authored task fragment may only say what the task runs. +FLOWX_OWNED_TASK_FIELDS = frozenset( + { + "task_key", + "depends_on", + "run_if", + "timeout_seconds", + "max_retries", + "min_retry_interval_millis", + "retry_on_timeout", + } +) + + +_COMPUTE_KEYS = ("environment_key", "job_cluster_key", "existing_cluster_id", "new_cluster") +_PAYLOAD_COMPUTE_KEYS = { + "spark_python_task": _COMPUTE_KEYS, + "python_wheel_task": _COMPUTE_KEYS, + "spark_jar_task": _COMPUTE_KEYS, + "dbt_task": _COMPUTE_KEYS, + "spark_submit_task": ("job_cluster_key", "new_cluster"), +} +_BUNDLE_JOB_CLUSTER_KEYS = (DEFAULT_JOB_CLUSTER_KEY, SINGLE_NODE_JOB_CLUSTER_KEY, MULTI_NODE_JOB_CLUSTER_KEY) + + +def _source_relative_path(raw_path: object) -> str: + """Return a safe path relative to the bundle's ``src`` directory.""" + relative_path = PurePosixPath(str(raw_path)) + windows_path = PureWindowsPath(str(raw_path)) + if ( + relative_path.is_absolute() + or windows_path.is_absolute() + # A drive or root without the other ("C:x.py", "\\tmp\\x.py") is not absolute on Windows but + # still resolves outside the bundle once joined to the output directory. + or windows_path.drive + or windows_path.root + or relative_path.as_posix() == "." + or ".." in relative_path.parts + or ".." in windows_path.parts + ): + raise ValueError(f"Agentic component file path {raw_path!r} must be relative to the bundle src directory") + return relative_path.as_posix() + + +def _check_task_payload(task_key: str, task: dict[str, object]) -> None: + """Require exactly one executable task payload (``notebook_task``, ``pipeline_task``, ...). + + Databricks names every task type ``_task``. A fragment with none packages fine but is + rejected at deploy time, and one with two leaves which runs undefined, so both fail here. + The same goes for a task that needs compute but names none it can run on, or names a job + cluster the bundle does not define. A ``for_each_task`` body is checked the same way. + """ + payloads = sorted(key for key in task if key.endswith("_task")) + if len(payloads) != 1: + found = ", ".join(payloads) if payloads else "none" + raise ValueError( + f"Agentic component {task_key!r} task must contain exactly one executable payload such as " + f"pipeline_task or notebook_task (found: {found})" + ) + payload = task[payloads[0]] + if not isinstance(payload, dict): + raise ValueError(f"Agentic component {task_key!r} task payload {payloads[0]} must be a mapping") + # flowx binds a cluster only for notebooks, and these task types cannot run without compute. + compute_keys = _PAYLOAD_COMPUTE_KEYS.get(payloads[0]) + if compute_keys and not any(key in task for key in compute_keys): + raise ValueError( + f"Agentic component {task_key!r} {payloads[0]} must name its compute with {', '.join(compute_keys)}" + ) + if "job_cluster_key" in task and task["job_cluster_key"] not in _BUNDLE_JOB_CLUSTER_KEYS: + raise ValueError( + f"Agentic component {task_key!r} job_cluster_key {task['job_cluster_key']!r} must be one of the " + f"bundle's job clusters: {', '.join(_BUNDLE_JOB_CLUSTER_KEYS)}" + ) + if payloads[0] == "for_each_task" and isinstance(payload.get("task"), dict): + _check_task_payload(task_key, payload["task"]) + + +def _check_resource(task_key: str, resource: object) -> None: + """Reject a resource that is not a mapping of a plain-identifier key to a definition mapping. + + Each resource is written to ``resources/.yml``, so a key with path characters + could land outside the ``resources`` directory or overwrite another bundle file. The definition + becomes the body of ``resources.pipelines.``, which the bundle expects to be a mapping. + """ + if not isinstance(resource, dict): + raise ValueError( + f"Agentic component {task_key!r} resource {resource!r} must be a mapping with resource_key and definition" + ) + resource_key = resource.get("resource_key") + if not isinstance(resource_key, str) or not resource_key or resource_key != normalize_task_key(resource_key): + raise ValueError( + f"Agentic component {task_key!r} resource key {resource_key!r} must be a plain identifier " + "of lowercase letters, digits, and single underscores" + ) + if not isinstance(resource.get("definition"), dict): + raise ValueError(f"Agentic component {task_key!r} resource {resource_key!r} definition must be a mapping") + + +def _check_environment(task_key: str, environment: object) -> None: + """Reject an environment that is not exactly an ``environment_key`` string and a ``spec`` mapping. + + The entry is written as given into a job's ``environments`` list, whose items take only those two fields. + """ + if ( + not isinstance(environment, dict) + or environment.keys() != {"environment_key", "spec"} + or not isinstance(environment["environment_key"], str) + or not environment["environment_key"] + or not isinstance(environment["spec"], dict) + ): + raise ValueError( + f"Agentic component {task_key!r} environment {environment!r} must be a mapping of a non-empty " + "environment_key string and a spec mapping" + ) + + +def _authored_file(task_key: str, file: object) -> DabNotebook: + """Turn one authored ``files`` entry into a file the bundle writer keeps below ``src``. + + A file carries exactly one of text ``content`` or base64 ``binary_content``, and it must + already be a string so the file is written byte for byte as authored. + """ + if not isinstance(file, dict) or not isinstance(file.get("path"), str): + raise ValueError(f"Agentic component {task_key!r} file {file!r} must be a mapping with a string path") + relative_path = _source_relative_path(file["path"]) + if ("content" in file) == ("binary_content" in file): + raise ValueError( + f"Agentic component {task_key!r} file {relative_path!r} must set content or binary_content, " + "exactly one of them" + ) + if "binary_content" in file: + encoded = file["binary_content"] + if not isinstance(encoded, str): + raise ValueError( + f"Agentic component {task_key!r} file {relative_path!r} binary_content must be a base64 string" + ) + try: + decoded = base64.b64decode(encoded, validate=True) + except ValueError as error: + raise ValueError( + f"Agentic component {task_key!r} file {relative_path!r} binary_content must be valid base64" + ) from error + return DabNotebook(relative_path=relative_path, binary_content=decoded, authored=True) + content = file["content"] + if not isinstance(content, str): + raise ValueError(f"Agentic component {task_key!r} file {relative_path!r} content must be a string") + return DabNotebook(relative_path=relative_path, content=content, authored=True) + + +def prepare(activity: AgenticComponentActivity, *, scope: str = "") -> PreparedActivity: + """Pass authored files, resources, environments, and the task payload to the bundle writer. + + A serverless task names its compute with ``environment_key``; the component declares + that environment in ``environments`` and the bundle writer adds it to the job that holds + the task. Every referenced ``environment_key`` must be declared by some component, and + two components may declare the same key only with an identical spec. A component cannot + declare a job cluster, so a ``job_cluster_key`` must name one the bundle writer defines: + ``default_cluster``, ``single_node_cluster`` or ``multi_node_cluster``. The bundle writer + adds that cluster to the job's ``job_clusters``. + + The task's key, dependencies, run condition, timeout, and retries always come from + the activity, never from the authored fragment, so an agent cannot rewire or re-time + a task behind flowx's back. A fragment that tries to set any of them is rejected. When + the component replaces a placeholder through ``merge_agentic``, the activity's own name, + key, dependencies, timeout, and retries are taken from that placeholder. + + flowx does not override or remove any other value the fragment sets, except in three + wiring passes that run on authored tasks exactly as on generated ones: + + - a ``{{tasks.X.values.Y}}`` reference to a task outside the same job is blanked, + because task values do not cross job boundaries. Only notebook ``base_parameters``, + ``run_job_task.job_parameters`` and ``condition_task`` operands are checked; + SETUP.md lists the blanked notebook parameters and neutralised conditions but not + blanked ``job_parameters``, and references in other payloads are left as they are; + - inside a ForEach that runs as its own job, notebook ``base_parameters`` and + ``condition_task`` operands are rewritten for that job: ``{{input.x}}`` becomes + ``{{job.parameters.x}}``, ADF expressions become job-parameter references, and + non-string values become strings; + - a ``run_job_task`` whose ``job_id`` is ``${resources.jobs.X.id}`` for a job outside + this bundle is pointed at an ``X_job_id`` bundle variable instead. + + Where the fragment leaves a value out, flowx adds only the plumbing the task needs: + + - the ForEach ``item`` parameter, only to a ``notebook_task`` (``base_parameters``) or + ``run_job_task`` (``job_parameters``) that is the ForEach's only child; any other + payload of an only child passes ``{{input}}`` itself. When the ForEach has several + children, or its only child is an IfCondition or Switch holding the component, the + children run as a ``_inner_tasks`` job where ``{{input}}`` does not resolve: + reference ``{{job.parameters.item}}`` in a notebook ``base_parameters``, a + ``run_job_task.job_parameters`` or a ``condition_task`` operand, which flowx forwards; + other payloads are not scanned, so they cannot receive the item on their own; + - the source activity's collapsed notifications, when the fragment sets neither + ``email_notifications`` nor ``webhook_notifications``; + - a job-cluster binding, only for a notebook task that names no compute and + either has a classic compute mode, ships libraries, or (without a serverless compute + mode) points at a workspace path outside the bundle. A bundle ``../src/`` notebook + without libraries runs on serverless, and no other task type is given a cluster. + + The task is marked ``_authored`` so the bundle writer skips the parameter clean-ups + meant for generated tasks; the marker is stripped before the job YAML is written. + + Raises: + ValueError: The authored task fragment sets a flowx-owned field, a file entry is + malformed or its path escapes the bundle's ``src`` directory, a resource is + malformed or its key is not a plain identifier, an environment is malformed, + a binary file is not valid base64, the task does not carry exactly one + executable payload mapping, a ``spark_python_task``, ``python_wheel_task``, + ``spark_jar_task`` or ``dbt_task`` names no compute, a ``spark_submit_task`` names + neither ``new_cluster`` nor a ``job_cluster_key`` (it runs only on a new cluster), or + a ``job_cluster_key`` is not one of the bundle's job clusters. + """ + del scope + owned_fields = sorted(FLOWX_OWNED_TASK_FIELDS & activity.task.keys()) + if owned_fields: + raise ValueError( + f"Agentic component {activity.task_key!r} task wiring sets flowx-owned field(s) " + f"{', '.join(owned_fields)}; set dependencies and run policy on the activity instead" + ) + _check_task_payload(activity.task_key, activity.task) + for resource in activity.resources: + _check_resource(activity.task_key, resource) + for environment in activity.environments: + _check_environment(activity.task_key, environment) + notebooks = [_authored_file(activity.task_key, file) for file in activity.files] + owned_task_fields = { + field: value for field, value in build_common_task_fields(activity).items() if field in FLOWX_OWNED_TASK_FIELDS + } + return PreparedActivity( + task={**owned_task_fields, **activity.task, "_authored": True}, + notebooks=notebooks, + pipeline_resources=list(activity.resources), + environments=list(activity.environments), + ) diff --git a/src/flowx/preparer/activity_preparers/for_each.py b/src/flowx/preparer/activity_preparers/for_each.py index a6570623..f61e18ac 100644 --- a/src/flowx/preparer/activity_preparers/for_each.py +++ b/src/flowx/preparer/activity_preparers/for_each.py @@ -115,20 +115,20 @@ def _resolve_for_each_inputs_with_bridge( def _inject_input_parameter(inner_task: dict) -> dict: - """Adds ``{{input}}`` as a base_parameter on the inner task. + """Adds ``{{input}}`` as the ``item`` parameter on the inner task unless it already passes one. Args: inner_task: The prepared inner task dict. Returns: - The task dict with ``item`` parameter injected. + The task dict with an ``item`` parameter. """ if "notebook_task" in inner_task: params = inner_task["notebook_task"].setdefault("base_parameters", {}) - params["item"] = "{{input}}" + params.setdefault("item", "{{input}}") elif "run_job_task" in inner_task: params = inner_task["run_job_task"].setdefault("job_parameters", {}) - params["item"] = "{{input}}" + params.setdefault("item", "{{input}}") return inner_task @@ -151,7 +151,7 @@ def prepare( Returns: A PreparedActivity with the for_each_task, plus any notebooks, secrets, - and inner_workflows from the child activities. + inner_workflows, pipeline resources and job environments from the child activities. """ task = build_common_task_fields(activity) concurrency = activity.concurrency if activity.concurrency is not None else 20 @@ -165,6 +165,8 @@ def prepare( all_secrets: list[SecretInstruction] = [] all_setup_tasks: list[SetupTask] = [] inner_workflows: list[PreparedWorkflow] = [] + pipeline_resources: list[dict[str, Any]] = [] + environments: list[dict[str, Any]] = [] extra_tasks: list[dict[str, Any]] = [] if inputs_bridge_task is not None: extra_tasks.append(inputs_bridge_task) @@ -175,6 +177,8 @@ def prepare( all_secrets.extend(inner_prepared.secrets) all_setup_tasks.extend(inner_prepared.setup_tasks) inner_workflows.extend(inner_prepared.inner_workflows) + pipeline_resources.extend(inner_prepared.pipeline_resources) + environments.extend(inner_prepared.environments) # CF-001: a single child with extra_tasks (IfCondition/Switch branch bodies) can't inline as # for_each_task.task (one task only); escalate to the sub-job path so the whole body survives. @@ -241,6 +245,8 @@ def prepare( all_secrets.extend(child_prepared.secrets) all_setup_tasks.extend(child_prepared.setup_tasks) inner_workflows.extend(child_prepared.inner_workflows) + pipeline_resources.extend(child_prepared.pipeline_resources) + environments.extend(child_prepared.environments) normalize_inner_task_params(inner_tasks) @@ -294,4 +300,6 @@ def prepare( secrets=all_secrets, setup_tasks=all_setup_tasks, inner_workflows=inner_workflows, + pipeline_resources=pipeline_resources, + environments=environments, ) diff --git a/src/flowx/preparer/activity_preparers/if_condition.py b/src/flowx/preparer/activity_preparers/if_condition.py index c50b4444..409c481d 100644 --- a/src/flowx/preparer/activity_preparers/if_condition.py +++ b/src/flowx/preparer/activity_preparers/if_condition.py @@ -118,6 +118,8 @@ def prepare(activity: IfConditionActivity, *, scope: str = "") -> PreparedActivi secrets=list(artifacts.secrets), setup_tasks=list(artifacts.setup_tasks), inner_workflows=list(artifacts.inner_workflows), + pipeline_resources=list(artifacts.pipeline_resources), + environments=list(artifacts.environments), ) diff --git a/src/flowx/preparer/activity_preparers/switch.py b/src/flowx/preparer/activity_preparers/switch.py index 5f5d4e0f..5d740dcf 100644 --- a/src/flowx/preparer/activity_preparers/switch.py +++ b/src/flowx/preparer/activity_preparers/switch.py @@ -127,6 +127,8 @@ def prepare(activity: SwitchActivity, *, scope: str = "") -> PreparedActivity: secrets=list(artifacts.secrets), setup_tasks=list(artifacts.setup_tasks), inner_workflows=list(artifacts.inner_workflows), + pipeline_resources=list(artifacts.pipeline_resources), + environments=list(artifacts.environments), ) # One condition task per case, chained via outcome="false" deps. Every case is named @@ -198,6 +200,8 @@ def prepare(activity: SwitchActivity, *, scope: str = "") -> PreparedActivity: secrets=list(artifacts.secrets), setup_tasks=list(artifacts.setup_tasks), inner_workflows=list(artifacts.inner_workflows), + pipeline_resources=list(artifacts.pipeline_resources), + environments=list(artifacts.environments), task_key_remap=remap, ) diff --git a/src/flowx/preparer/workflow_preparer.py b/src/flowx/preparer/workflow_preparer.py index 3ab1df13..1ed5528e 100644 --- a/src/flowx/preparer/workflow_preparer.py +++ b/src/flowx/preparer/workflow_preparer.py @@ -9,6 +9,7 @@ from flowx.models.dab import DabNotebook, ParameterApproximation, SecretInstruction, SetupTask from flowx.models.ir import ( Activity, + AgenticComponentActivity, AppendVariableActivity, CopyActivity, DbtFactoryActivity, @@ -33,6 +34,8 @@ WebActivity, ) +_TASK_NOTIFICATION_KEYS = ("email_notifications", "webhook_notifications") + @dataclass(slots=True, kw_only=True) class PreparedActivity: @@ -48,6 +51,8 @@ class PreparedActivity: task_key_remap: dict[str, str] = field(default_factory=dict) # Lakeflow pipeline resources emitted under resources/.yml ({resource_key, definition}). pipeline_resources: list[dict[str, Any]] = field(default_factory=list) + # Job environments ({environment_key, spec}) written into each job whose tasks reference them. + environments: list[dict[str, Any]] = field(default_factory=list) parameter_approximations: list[ParameterApproximation] = field(default_factory=list) @@ -64,6 +69,7 @@ class PreparedWorkflow: parameters: list[dict[str, Any]] = field(default_factory=list) cluster_hints: list[dict[str, Any]] = field(default_factory=list) pipeline_resources: list[dict[str, Any]] = field(default_factory=list) + environments: list[dict[str, Any]] = field(default_factory=list) parameter_approximations: list[ParameterApproximation] = field(default_factory=list) # C-10 (SCHED-001): serialised schedule / trigger spec the bundler # renders as ``schedule:`` / ``trigger:`` on the emitted DAB job. @@ -142,6 +148,7 @@ def prepare_activity( ) -> PreparedActivity: """Dispatches to the appropriate activity preparer based on activity type.""" from flowx.preparer.activity_preparers import ( + agentic_component, append_variable, copy, databricks_job, @@ -164,6 +171,7 @@ def prepare_activity( ) dispatch: dict[type, Any] = { + AgenticComponentActivity: agentic_component.prepare, NotebookActivity: notebook.prepare, SparkJarActivity: spark_jar.prepare, SparkPythonActivity: spark_python.prepare, @@ -207,7 +215,7 @@ def prepare_activity( prepared.task = _stamp_compute_mode(prepared.task, activity.compute_mode) # Wire any notification spec the adapter stamped onto this task (generic across task types, not just Copy). - if activity.notifications: + if activity.notifications and not any(key in prepared.task for key in _TASK_NOTIFICATION_KEYS): from flowx.preparer.notifications import resolve_task_notifications notification_keys, notification_setup = resolve_task_notifications(activity.notifications) @@ -301,6 +309,7 @@ class PreparedArtifacts: setup_tasks: tuple[SetupTask, ...] = () inner_workflows: tuple[PreparedWorkflow, ...] = () pipeline_resources: tuple[dict[str, Any], ...] = () + environments: tuple[dict[str, Any], ...] = () parameter_approximations: tuple[ParameterApproximation, ...] = () @@ -315,6 +324,7 @@ def merge_prepared_artifacts( setup_tasks=artifacts.setup_tasks + tuple(prepared.setup_tasks), inner_workflows=artifacts.inner_workflows + tuple(prepared.inner_workflows), pipeline_resources=artifacts.pipeline_resources + tuple(prepared.pipeline_resources), + environments=artifacts.environments + tuple(prepared.environments), parameter_approximations=artifacts.parameter_approximations + tuple(prepared.parameter_approximations), ) @@ -426,6 +436,7 @@ def prepare_workflow(pipeline: Pipeline) -> PreparedWorkflow: inner_workflows=list(artifacts.inner_workflows), cluster_hints=cluster_hints, pipeline_resources=list(artifacts.pipeline_resources), + environments=list(artifacts.environments), parameter_approximations=list(artifacts.parameter_approximations), schedule=pipeline.schedule, bundle_variables=dict(pipeline.bundle_variables), diff --git a/src/flowx/validate/bundle_invariants.py b/src/flowx/validate/bundle_invariants.py index 2ed275f6..9e689a99 100644 --- a/src/flowx/validate/bundle_invariants.py +++ b/src/flowx/validate/bundle_invariants.py @@ -25,6 +25,7 @@ _ANCHOR_RE = re.compile(r"[&*]id\d+\b") _JOB_PARAM_REF_RE = re.compile(r"\{\{\s*job\.parameters\.([A-Za-z0-9_]+)\s*\}\}") _JOB_RESOURCE_ID_RE = re.compile(r"\$\{resources\.jobs\.([^.}]+)\.id\}") +_PIPELINE_RESOURCE_ID_RE = re.compile(r"\$\{resources\.pipelines\.([^.}]+)\.id\}") _PYDABS_JOB_RE = re.compile(r"resources\.add_job\(\s*['\"]([^'\"]+)['\"]") @@ -246,10 +247,15 @@ def check_bundle_dir(bundle_dir: Path) -> BundleInvariantResult: documents.append((path, document)) known_jobs: set[str] = set() + known_pipelines: set[str] = set() for _path, document in documents: - jobs = (document.get("resources") or {}).get("jobs") or {} + document_resources = document.get("resources") or {} + jobs = document_resources.get("jobs") or {} if isinstance(jobs, dict): known_jobs.update(str(job_key) for job_key in jobs) + pipelines = document_resources.get("pipelines") or {} + if isinstance(pipelines, dict): + known_pipelines.update(str(pipeline_key) for pipeline_key in pipelines) python_resources = (document.get("python") or {}).get("resources") or [] for resource in python_resources: if not isinstance(resource, str): @@ -267,7 +273,39 @@ def check_bundle_dir(bundle_dir: Path) -> BundleInvariantResult: for job_key, job in jobs.items(): if not isinstance(job, dict): continue + job_environment_keys = { + environment.get("environment_key") + for environment in job.get("environments") or [] + if isinstance(environment, dict) + } for task in _iter_tasks(job.get("tasks") or []): + if "environment_key" in task and task["environment_key"] not in job_environment_keys: + findings.append( + BundleFinding( + code="dangling_environment_reference", + location=f"{path.name}, job '{job_key}', task '{task.get('task_key', '')}'", + message=( + f"Task environment_key '{task['environment_key']}' is not declared in the job's " + "environments." + ), + ) + ) + pipeline_task = task.get("pipeline_task") or {} + pipeline_id = pipeline_task.get("pipeline_id") if isinstance(pipeline_task, dict) else None + pipeline_match = ( + _PIPELINE_RESOURCE_ID_RE.fullmatch(pipeline_id) if isinstance(pipeline_id, str) else None + ) + if pipeline_match is not None and pipeline_match.group(1) not in known_pipelines: + findings.append( + BundleFinding( + code="dangling_pipeline_reference", + location=f"{path.name}, job '{job_key}', task '{task.get('task_key', '')}'", + message=( + f"pipeline_task references bundle pipeline '{pipeline_match.group(1)}', which is not " + "declared in static resource YAML." + ), + ) + ) run_job = task.get("run_job_task") or {} job_id = run_job.get("job_id") if isinstance(run_job, dict) else None match = _JOB_RESOURCE_ID_RE.fullmatch(job_id) if isinstance(job_id, str) else None diff --git a/tests/unit/test_agentic_component.py b/tests/unit/test_agentic_component.py new file mode 100644 index 00000000..386e3341 --- /dev/null +++ b/tests/unit/test_agentic_component.py @@ -0,0 +1,1004 @@ +"""Tests for generic agent-authored bundle components.""" + +from __future__ import annotations + +import base64 +import json + +import pytest +import yaml + +from flowx.bundler.dab_writer import _combine_airflow_workflows, pipeline_dict_to_ir, write_bundle +from flowx.ir_serde import activity_to_dict, merge_agentic_results, pipeline_to_dict +from flowx.models.dab import SetupTask +from flowx.models.ir import ( + AgenticComponentActivity, + Dependency, + ForEachActivity, + IfConditionActivity, + Pipeline, + PlaceholderActivity, + SwitchActivity, + SwitchCase, + WaitActivity, +) +from flowx.preparer.workflow_preparer import prepare_workflow +from flowx.validate.bundle_invariants import check_bundle_dir + +SOURCE_DEFINITION = {"type": "ExecuteDataFlow", "typeProperties": {"dataflow": "orders"}} +FILES = [ + {"path": "pipelines/orders.py", "content": "from pyspark import pipelines as dp\n"}, + {"path": "libraries/orders.whl", "binary_content": "UEsDBAoAAAAA"}, +] +PIPELINE_DEFINITION = { + "name": "orders_ingestion", + "catalog": "${var.catalog}", + "target": "${var.schema}", + "ingestion_definition": { + "connection_name": "flowx_orders_connection", + "objects": [ + { + "table": { + "source_catalog": "sales", + "source_schema": "dbo", + "source_table": "orders", + "destination_catalog": "${var.catalog}", + "destination_schema": "${var.schema}", + "destination_table": "orders", + } + } + ], + }, +} +RESOURCES = [{"resource_key": "orders_ingestion", "definition": PIPELINE_DEFINITION}] +TASK = {"pipeline_task": {"pipeline_id": "${resources.pipelines.orders_ingestion.id}"}} +SERVERLESS_ENVIRONMENT = { + "environment_key": "serverless", + "spec": {"environment_version": "2", "dependencies": ["requests==2.32.3"]}, +} +WHEEL_ENVIRONMENT = { + "environment_key": "serverless", + "spec": {"environment_version": "2", "dependencies": ["../src/libraries/orders.whl"]}, +} + + +def _activity() -> AgenticComponentActivity: + return AgenticComponentActivity( + name="Ingest orders", + task_key="ingest_orders", + files=FILES, + resources=RESOURCES, + task=TASK, + raw_definition=SOURCE_DEFINITION, + ) + + +def test_agentic_component_round_trips_through_serialized_ir(): + serialized = json.loads(json.dumps(activity_to_dict(_activity()))) + pipeline, _ = pipeline_dict_to_ir({"name": "orders", "tasks": [serialized]}) + restored = pipeline.tasks[0] + + assert type(restored).__name__ == "AgenticComponentActivity" + assert activity_to_dict(restored) == serialized + assert serialized == { + "name": "Ingest orders", + "task_key": "ingest_orders", + "type": "AgenticComponentActivity", + "files": FILES, + "resources": RESOURCES, + "task": TASK, + "raw_definition": SOURCE_DEFINITION, + } + + +def test_agentic_component_packages_authored_files_resource_and_pipeline_task(tmp_path): + serialized_pipeline = json.loads(json.dumps(pipeline_to_dict(Pipeline(name="orders", tasks=[_activity()])))) + restored_pipeline, _ = pipeline_dict_to_ir(serialized_pipeline) + + write_bundle(prepare_workflow(restored_pipeline), tmp_path) + + assert (tmp_path / "src" / "pipelines" / "orders.py").read_text(encoding="utf-8") == FILES[0]["content"] + assert (tmp_path / "src" / "libraries" / "orders.whl").read_bytes() == base64.b64decode(FILES[1]["binary_content"]) + pipeline_resource = yaml.safe_load((tmp_path / "resources" / "orders_ingestion.yml").read_text(encoding="utf-8")) + assert pipeline_resource == {"resources": {"pipelines": {"orders_ingestion": PIPELINE_DEFINITION}}} + + job_resource = yaml.safe_load((tmp_path / "resources" / "orders.yml").read_text(encoding="utf-8")) + assert job_resource["resources"]["jobs"]["orders"]["tasks"] == [ + {"task_key": "ingest_orders", "pipeline_task": TASK["pipeline_task"]} + ] + assert check_bundle_dir(tmp_path).ok + + +def test_agentic_component_notebook_files_with_reserved_paths_stay_under_src(tmp_path): + activity = AgenticComponentActivity( + name="Custom notebook", + task_key="custom_notebook", + files=[ + {"path": "resources/custom.py", "content": "print('custom')\n"}, + {"path": "pyproject.toml", "content": "[project]\nname = 'custom'\n"}, + ], + task={"notebook_task": {"notebook_path": "../src/resources/custom.py"}}, + ) + + write_bundle(prepare_workflow(Pipeline(name="custom", tasks=[activity])), tmp_path) + + assert (tmp_path / "src" / "resources" / "custom.py").read_text(encoding="utf-8") == "print('custom')\n" + assert (tmp_path / "src" / "pyproject.toml").read_text(encoding="utf-8") == "[project]\nname = 'custom'\n" + assert not (tmp_path / "resources" / "custom.py").exists() + assert not (tmp_path / "pyproject.toml").exists() + job_resource = yaml.safe_load((tmp_path / "resources" / "custom.yml").read_text(encoding="utf-8")) + assert job_resource["resources"]["jobs"]["custom"]["tasks"][0]["notebook_task"]["notebook_path"] == ( + "../src/resources/custom.py" + ) + + +def test_agentic_component_notebook_path_tracks_shared_workflow_namespacing(tmp_path): + custom = AgenticComponentActivity( + name="Custom notebook", + task_key="custom_notebook", + files=[{"path": "resources/custom.py", "content": "print('custom')\n"}], + task={"notebook_task": {"notebook_path": "../src/resources/custom.py"}}, + ) + other = AgenticComponentActivity( + name="Other notebook", + task_key="other_notebook", + files=[{"path": "notebooks/other.py", "content": "print('other')\n"}], + task={"notebook_task": {"notebook_path": "../src/notebooks/other.py"}}, + ) + combined = _combine_airflow_workflows( + [ + prepare_workflow(Pipeline(name="custom", tasks=[custom])), + prepare_workflow(Pipeline(name="other", tasks=[other])), + ] + ) + + write_bundle(combined, tmp_path) + + assert (tmp_path / "src" / "custom" / "resources" / "custom.py").exists() + custom_job = yaml.safe_load((tmp_path / "resources" / "custom.yml").read_text(encoding="utf-8")) + assert custom_job["resources"]["jobs"]["custom"]["tasks"][0]["notebook_task"]["notebook_path"] == ( + "../src/custom/resources/custom.py" + ) + assert (tmp_path / "src" / "other" / "notebooks" / "other.py").exists() + other_job = yaml.safe_load((tmp_path / "resources" / "other.yml").read_text(encoding="utf-8")) + assert other_job["resources"]["jobs"]["other"]["tasks"][0]["notebook_task"]["notebook_path"] == ( + "../src/other/notebooks/other.py" + ) + + +@pytest.mark.parametrize( + "path", + ["../outside.py", "/tmp/outside.py", "..\\outside.py", "C:\\tmp\\outside.py"], +) +def test_agentic_component_rejects_file_paths_that_escape_src(path): + activity = AgenticComponentActivity( + name="Unsafe file", + task_key="unsafe_file", + files=[{"path": path, "content": "unsafe\n"}], + task={"notebook_task": {"notebook_path": "../src/safe.py"}}, + ) + + with pytest.raises(ValueError, match="relative to the bundle src directory"): + prepare_workflow(Pipeline(name="unsafe", tasks=[activity])) + + +@pytest.mark.parametrize( + "owned_field", + [ + "task_key", + "depends_on", + "run_if", + "timeout_seconds", + "max_retries", + "min_retry_interval_millis", + "retry_on_timeout", + ], +) +def test_agentic_component_rejects_task_wiring_that_sets_flowx_owned_fields(owned_field): + activity = AgenticComponentActivity( + name="Rewired", + task_key="rewired", + resources=RESOURCES, + task={**TASK, owned_field: "authored"}, + ) + + with pytest.raises(ValueError, match="flowx-owned"): + prepare_workflow(Pipeline(name="rewired", tasks=[activity])) + + +def test_agentic_component_task_identity_dependencies_and_policy_come_from_the_activity(tmp_path): + upstream = AgenticComponentActivity(name="Upstream", task_key="upstream", resources=RESOURCES, task=TASK) + downstream = AgenticComponentActivity( + name="Downstream", + task_key="downstream", + depends_on=[Dependency(task_key="upstream", outcome="Succeeded")], + timeout_seconds=600, + max_retries=2, + task={"notebook_task": {"notebook_path": "../src/notebooks/downstream.py"}}, + files=[{"path": "notebooks/downstream.py", "content": "print('downstream')\n"}], + ) + + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[upstream, downstream])), tmp_path) + + job_resource = yaml.safe_load((tmp_path / "resources" / "orders.yml").read_text(encoding="utf-8")) + tasks = {task["task_key"]: task for task in job_resource["resources"]["jobs"]["orders"]["tasks"]} + assert tasks["downstream"]["depends_on"] == [{"task_key": "upstream"}] + assert tasks["downstream"]["timeout_seconds"] == 600 + assert tasks["downstream"]["max_retries"] == 2 + assert "notebook_task" in tasks["downstream"] + assert check_bundle_dir(tmp_path).ok + + +def test_merged_agentic_component_takes_identity_dependencies_and_policy_from_the_placeholder(tmp_path): + """The agent's own key, wiring, and run policy on a merged component are replaced by the placeholder's.""" + report = tmp_path / "translation_report.json" + placeholders = Pipeline( + name="orders", + tasks=[ + WaitActivity(name="Extract", task_key="extract", wait_time_seconds=5), + PlaceholderActivity( + name="Transform", + task_key="transform", + original_type="ExecuteDataFlow", + depends_on=[Dependency(task_key="extract", outcome="Succeeded")], + timeout_seconds=3600, + max_retries=2, + min_retry_interval_millis=60000, + ), + PlaceholderActivity(name="Load", task_key="load", original_type="Custom"), + ], + ) + report.write_text(json.dumps(pipeline_to_dict(placeholders)), encoding="utf-8") + results = tmp_path / "agentic_results" + results.mkdir() + transform = { + "type": "AgenticComponentActivity", + "name": "Renamed", + "task_key": "rewired", + "depends_on": [], + "timeout_seconds": 5, + "resources": RESOURCES, + "task": TASK, + } + load = { + "type": "AgenticComponentActivity", + "depends_on": [{"task_key": "extract", "outcome": "Succeeded"}], + "max_retries": 9, + "files": [{"path": "jobs/load.py", "content": "print('load')\n"}], + "task": {"spark_python_task": {"python_file": "../src/jobs/load.py"}, "existing_cluster_id": "0101-abc"}, + } + (results / "transform.json").write_text( + json.dumps({"activity_name": "Transform", "task": transform}), encoding="utf-8" + ) + (results / "load.json").write_text(json.dumps({"activity_name": "Load", "task": load}), encoding="utf-8") + + assert merge_agentic_results(report, results) == (2, 0) + merged_pipeline, _ = pipeline_dict_to_ir(json.loads(report.read_text(encoding="utf-8"))) + write_bundle(prepare_workflow(merged_pipeline), tmp_path / "bundle") + + job_resource = yaml.safe_load((tmp_path / "bundle" / "resources" / "orders.yml").read_text(encoding="utf-8")) + tasks = {task["task_key"]: task for task in job_resource["resources"]["jobs"]["orders"]["tasks"]} + assert tasks["transform"] == { + "task_key": "transform", + "depends_on": [{"task_key": "extract"}], + "timeout_seconds": 3600, + "retry_on_timeout": True, + "max_retries": 2, + "min_retry_interval_millis": 60000, + "pipeline_task": TASK["pipeline_task"], + } + assert tasks["load"] == { + "task_key": "load", + "spark_python_task": {"python_file": "../src/jobs/load.py"}, + "existing_cluster_id": "0101-abc", + } + assert [task.name for task in merged_pipeline.tasks] == ["Extract", "Transform", "Load"] + + +@pytest.mark.parametrize("resource_key", ["../databricks", "../../outside", "nested/pipeline", "nested\\pipeline", ""]) +def test_agentic_component_rejects_resource_keys_that_are_not_plain_identifiers(resource_key): + activity = AgenticComponentActivity( + name="Unsafe resource", + task_key="unsafe_resource", + resources=[{"resource_key": resource_key, "definition": PIPELINE_DEFINITION}], + task=TASK, + ) + + with pytest.raises(ValueError, match="must be a plain identifier"): + prepare_workflow(Pipeline(name="unsafe", tasks=[activity])) + + +def test_agentic_component_resource_key_matching_the_job_key_fails_instead_of_overwriting_the_job(tmp_path): + activity = AgenticComponentActivity(name="Ingest orders", task_key="ingest_orders", resources=RESOURCES, task=TASK) + + with pytest.raises(ValueError, match="matches a job resource key"): + write_bundle(prepare_workflow(Pipeline(name="orders_ingestion", tasks=[activity])), tmp_path) + + assert not (tmp_path / "resources" / "orders_ingestion.yml").exists() + + +def test_agentic_component_resource_key_matching_a_dbt_factory_job_key_fails(tmp_path): + activity = AgenticComponentActivity( + name="Ingest orders", + task_key="ingest_orders", + resources=[{"resource_key": "orders_dbt", "definition": PIPELINE_DEFINITION}], + task={"pipeline_task": {"pipeline_id": "${resources.pipelines.orders_dbt.id}"}}, + ) + workflow = prepare_workflow(Pipeline(name="orders", tasks=[activity])) + workflow.setup_tasks.append( + SetupTask( + type="pydabs_dbt_factory", + config={"hook_module": "resources.orders_dbt_job", "job_key": "orders_dbt"}, + ) + ) + + with pytest.raises(ValueError, match="matches a job resource key"): + write_bundle(workflow, tmp_path) + + +def test_agentic_components_declaring_an_identical_pipeline_resource_write_it_once(tmp_path): + first = AgenticComponentActivity(name="First", task_key="first", resources=RESOURCES, task=TASK) + second = AgenticComponentActivity(name="Second", task_key="second", resources=RESOURCES, task=TASK) + serialized_pipeline = json.loads(json.dumps(pipeline_to_dict(Pipeline(name="orders", tasks=[first, second])))) + restored_pipeline, _ = pipeline_dict_to_ir(serialized_pipeline) + + created_files = write_bundle(prepare_workflow(restored_pipeline), tmp_path) + + resource_path = (tmp_path / "resources" / "orders_ingestion.yml").resolve() + assert created_files.count(resource_path) == 1 + pipeline_resource = yaml.safe_load(resource_path.read_text(encoding="utf-8")) + assert pipeline_resource == {"resources": {"pipelines": {"orders_ingestion": PIPELINE_DEFINITION}}} + job_resource = yaml.safe_load((tmp_path / "resources" / "orders.yml").read_text(encoding="utf-8")) + assert [task["pipeline_task"] for task in job_resource["resources"]["jobs"]["orders"]["tasks"]] == [ + TASK["pipeline_task"], + TASK["pipeline_task"], + ] + assert check_bundle_dir(tmp_path).ok + + +def test_agentic_components_declaring_one_pipeline_definition_with_different_extra_fields_write_it_once(tmp_path): + first = AgenticComponentActivity( + name="First", + task_key="first", + resources=[{**RESOURCES[0], "comment": "for first"}], + task=TASK, + ) + second = AgenticComponentActivity( + name="Second", + task_key="second", + resources=[{**RESOURCES[0], "comment": "for second"}], + task=TASK, + ) + serialized_pipeline = json.loads(json.dumps(pipeline_to_dict(Pipeline(name="orders", tasks=[first, second])))) + restored_pipeline, _ = pipeline_dict_to_ir(serialized_pipeline) + + created_files = write_bundle(prepare_workflow(restored_pipeline), tmp_path) + + resource_path = (tmp_path / "resources" / "orders_ingestion.yml").resolve() + assert created_files.count(resource_path) == 1 + pipeline_resource = yaml.safe_load(resource_path.read_text(encoding="utf-8")) + assert pipeline_resource == {"resources": {"pipelines": {"orders_ingestion": PIPELINE_DEFINITION}}} + + +def test_agentic_components_sharing_a_resource_key_with_different_definitions_fail(tmp_path): + first = AgenticComponentActivity(name="First", task_key="first", resources=RESOURCES, task=TASK) + second = AgenticComponentActivity( + name="Second", + task_key="second", + resources=[{"resource_key": "orders_ingestion", "definition": {**PIPELINE_DEFINITION, "name": "other"}}], + task=TASK, + ) + + with pytest.raises(ValueError, match="more than one pipeline resource with different definitions"): + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[first, second])), tmp_path) + assert not (tmp_path / "databricks.yml").exists() + + +@pytest.mark.parametrize( + ("files", "resources", "task"), + [ + ([{"path": "notebooks/a.py", "content": None}], [], TASK), + ([{"path": "notebooks/a.py", "content": ["line 1", "line 2"]}], [], TASK), + ([{"path": "libraries/a.whl", "binary_content": None}], [], TASK), + ([{"path": "libraries/a.whl", "content": "text", "binary_content": "UEsDBAoAAAAA"}], [], TASK), + (["notebooks/a.py"], [], TASK), + ([{"content": "print('no path')\n"}], [], TASK), + ([{"path": None, "content": "print('null path')\n"}], [], TASK), + ([], ["orders_ingestion"], TASK), + ([], [{"resource_key": "orders_ingestion"}], TASK), + ([], [{"resource_key": "orders_ingestion", "definition": "orders"}], TASK), + ([], [{"resource_key": "orders_ingestion", "definition": ["orders"]}], TASK), + ([], [], {"notebook_task": "../src/notebooks/a.py"}), + ], +) +def test_agentic_component_rejects_malformed_authored_payloads(files, resources, task): + activity = AgenticComponentActivity(name="Bad", task_key="bad", files=files, resources=resources, task=task) + + with pytest.raises(ValueError, match="Agentic component 'bad'"): + prepare_workflow(Pipeline(name="orders", tasks=[activity])) + + +def test_agentic_component_task_fragment_keeps_its_own_description_and_compute(tmp_path): + """Only the flowx-owned fields come from the activity; the fragment's compute is not overridden.""" + activity = AgenticComponentActivity( + name="Run job", + task_key="run_job", + description="Copied from the source activity", + existing_cluster_id="0101-123456-abcdefgh", + files=[{"path": "jobs/run.py", "content": "print('run')\n"}], + task={"spark_python_task": {"python_file": "../src/jobs/run.py"}, "job_cluster_key": "default_cluster"}, + ) + + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[activity])), tmp_path) + + job_resource = yaml.safe_load((tmp_path / "resources" / "orders.yml").read_text(encoding="utf-8")) + assert job_resource["resources"]["jobs"]["orders"]["tasks"] == [ + { + "task_key": "run_job", + "spark_python_task": {"python_file": "../src/jobs/run.py"}, + "job_cluster_key": "default_cluster", + } + ] + + +def test_agentic_component_rejects_binary_content_that_is_not_plain_base64(): + activity = AgenticComponentActivity( + name="Wheel", + task_key="wheel", + files=[{"path": "libraries/orders.whl", "binary_content": "data:application/zip;base64,UEsDBAoAAAAA"}], + task={"notebook_task": {"notebook_path": "../src/notebooks/orders.py"}}, + ) + + with pytest.raises( + ValueError, match="Agentic component 'wheel' file 'libraries/orders.whl' binary_content must be valid base64" + ): + prepare_workflow(Pipeline(name="orders", tasks=[activity])) + + +@pytest.mark.parametrize("path", ["\\tmp\\outside.py", "C:outside.py", "\\\\server\\share\\outside.py"]) +def test_agentic_component_rejects_windows_rooted_or_drive_relative_paths(path): + """A path with a Windows root or drive but no full absolute form still escapes the bundle.""" + activity = AgenticComponentActivity( + name="Unsafe file", + task_key="unsafe_file", + files=[{"path": path, "content": "unsafe\n"}], + task={"notebook_task": {"notebook_path": "../src/safe.py"}}, + ) + + with pytest.raises(ValueError, match="relative to the bundle src directory"): + prepare_workflow(Pipeline(name="unsafe", tasks=[activity])) + + +@pytest.mark.parametrize( + "task", + [ + {}, + {"notebook_task": {"notebook_path": "../src/a.py"}, "pipeline_task": {"pipeline_id": "x"}}, + ], +) +def test_agentic_component_requires_exactly_one_executable_payload(task): + activity = AgenticComponentActivity(name="No payload", task_key="no_payload", task=task) + + with pytest.raises(ValueError, match="exactly one executable payload"): + prepare_workflow(Pipeline(name="orders", tasks=[activity])) + + +def test_agentic_components_authoring_different_content_at_one_path_fail(tmp_path): + """The writer would rename the second file, leaving its task running the first component's code.""" + first = AgenticComponentActivity( + name="First", + task_key="first", + files=[{"path": "jobs/run.py", "content": "print('first')\n"}], + task={"spark_python_task": {"python_file": "../src/jobs/run.py"}, "existing_cluster_id": "0101-abc"}, + ) + second = AgenticComponentActivity( + name="Second", + task_key="second", + files=[{"path": "jobs/run.py", "content": "print('second')\n"}], + task={"spark_python_task": {"python_file": "../src/jobs/run.py"}, "existing_cluster_id": "0101-abc"}, + ) + + with pytest.raises(ValueError, match="'jobs/run.py' shares its path"): + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[first, second])), tmp_path) + assert not (tmp_path / "databricks.yml").exists() + assert not (tmp_path / "resources").exists() + assert not (tmp_path / "src").exists() + + +def test_agentic_components_sharing_identical_file_content_still_package(tmp_path): + shared = [{"path": "jobs/common.py", "content": "print('shared')\n"}] + first = AgenticComponentActivity( + name="First", + task_key="first", + files=shared, + task={"spark_python_task": {"python_file": "../src/jobs/common.py"}, "existing_cluster_id": "0101-abc"}, + ) + second = AgenticComponentActivity( + name="Second", + task_key="second", + files=shared, + task={"spark_python_task": {"python_file": "../src/jobs/common.py"}, "existing_cluster_id": "0101-abc"}, + ) + + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[first, second])), tmp_path) + + assert (tmp_path / "src" / "jobs" / "common.py").read_text(encoding="utf-8") == "print('shared')\n" + + +def test_agentic_component_file_at_a_generated_notebook_path_fails(tmp_path): + """The writer would rename one of the two files, leaving one task running the other's code.""" + wait = WaitActivity(name="Pause", task_key="pause", wait_time_seconds=5) + custom = AgenticComponentActivity( + name="Custom", + task_key="custom", + files=[{"path": "notebooks/pause.py", "content": "print('custom')\n"}], + task={"notebook_task": {"notebook_path": "../src/notebooks/pause.py"}}, + ) + + with pytest.raises(ValueError, match="'notebooks/pause.py'"): + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[wait, custom])), tmp_path) + assert not (tmp_path / "src").exists() + + +def test_agentic_component_file_at_a_generated_setup_notebook_path_fails(tmp_path): + """Setup notebooks are written after authored files and would silently replace them.""" + custom = AgenticComponentActivity( + name="Custom", + task_key="custom", + files=[{"path": "setup/create_volumes.py", "content": "print('custom')\n"}], + task={"notebook_task": {"notebook_path": "../src/setup/create_volumes.py"}}, + ) + workflow = prepare_workflow(Pipeline(name="orders", tasks=[custom])) + workflow.setup_tasks.append(SetupTask(type="volume", config={"volume_name": "landing"})) + + with pytest.raises(ValueError, match="'setup/create_volumes.py'"): + write_bundle(workflow, tmp_path) + assert not (tmp_path / "src").exists() + + +def _inside_if_condition(activity: AgenticComponentActivity) -> IfConditionActivity: + return IfConditionActivity( + name="Check", task_key="check", op="EQUAL_TO", left="1", right="1", if_true_activities=[activity] + ) + + +def _inside_switch(activity: AgenticComponentActivity) -> SwitchActivity: + return SwitchActivity( + name="Route", + task_key="route", + on_expression="orders", + cases=[SwitchCase(value="orders", activities=[activity])], + ) + + +def _inside_for_each(activity: AgenticComponentActivity) -> ForEachActivity: + return ForEachActivity(name="Loop", task_key="loop", items_expression='["a"]', inner_activities=[activity]) + + +def _inside_for_each_with_siblings(activity: AgenticComponentActivity) -> ForEachActivity: + sibling = WaitActivity(name="Pause", task_key="pause", wait_time_seconds=5) + return ForEachActivity(name="Loop", task_key="loop", items_expression='["a"]', inner_activities=[activity, sibling]) + + +@pytest.mark.parametrize( + "wrap", [_inside_if_condition, _inside_switch, _inside_for_each, _inside_for_each_with_siblings] +) +def test_agentic_component_pipeline_resource_survives_control_flow(tmp_path, wrap): + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[wrap(_activity())])), tmp_path) + + pipeline_resource = yaml.safe_load((tmp_path / "resources" / "orders_ingestion.yml").read_text(encoding="utf-8")) + assert pipeline_resource == {"resources": {"pipelines": {"orders_ingestion": PIPELINE_DEFINITION}}} + assert check_bundle_dir(tmp_path).ok + + +def _notebook_tasks_by_path(bundle_dir) -> dict[str, dict]: + """Collect every emitted task with a ``notebook_task`` across the bundle's job resources, keyed by notebook path.""" + found: dict[str, dict] = {} + + def visit(value) -> None: + if isinstance(value, dict): + notebook_task = value.get("notebook_task") + if isinstance(notebook_task, dict): + found[notebook_task["notebook_path"]] = value + for item in value.values(): + visit(item) + elif isinstance(value, list): + for item in value: + visit(item) + + for resource_path in sorted((bundle_dir / "resources").glob("*.yml")): + visit(yaml.safe_load(resource_path.read_text(encoding="utf-8"))) + return found + + +@pytest.mark.parametrize( + "wrap", + [lambda activity: activity, _inside_for_each, _inside_for_each_with_siblings], + ids=["top_level", "for_each", "for_each_with_siblings"], +) +def test_agentic_component_notebook_widget_defaults_are_not_overridden(tmp_path, wrap): + """The authored notebook's own widget default must apply, not an injected empty base parameter.""" + report = AgenticComponentActivity( + name="Report", + task_key="report", + files=[ + { + "path": "notebooks/report.py", + "content": ( + 'dbutils.widgets.text("lookback_days", "7")\n' + 'lookback_days = int(dbutils.widgets.get("lookback_days"))\n' + ), + } + ], + task={"notebook_task": {"notebook_path": "../src/notebooks/report.py"}}, + ) + + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[wrap(report)])), tmp_path) + + task = _notebook_tasks_by_path(tmp_path)["../src/notebooks/report.py"] + assert "lookback_days" not in task["notebook_task"].get("base_parameters", {}) + assert "_authored" not in task + + +@pytest.mark.parametrize( + ("fragment", "parameters_field"), + [ + ( + { + "notebook_task": { + "notebook_path": "../src/notebooks/load.py", + "base_parameters": {"item": "{{input.table_name}}"}, + } + }, + ("notebook_task", "base_parameters"), + ), + ( + {"run_job_task": {"job_id": 123, "job_parameters": {"item": "{{input.table_name}}"}}}, + ("run_job_task", "job_parameters"), + ), + ], + ids=["notebook_task", "run_job_task"], +) +def test_agentic_component_inside_for_each_keeps_its_own_item_parameter(tmp_path, fragment, parameters_field): + load = AgenticComponentActivity( + name="Load", + task_key="load", + files=[{"path": "notebooks/load.py", "content": "print('load')\n"}], + task=fragment, + ) + + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[_inside_for_each(load)])), tmp_path) + + job_resource = yaml.safe_load((tmp_path / "resources" / "orders.yml").read_text(encoding="utf-8")) + body = job_resource["resources"]["jobs"]["orders"]["tasks"][0]["for_each_task"]["task"] + payload, parameters = parameters_field + assert body[payload][parameters] == {"item": "{{input.table_name}}"} + assert "_authored" not in body + + +def test_agentic_component_notifications_are_added_only_when_the_fragment_sets_none(tmp_path): + collapsed = {"destination": "email", "events": ["on_failure"], "args": {"addresses": ["flowx@example.com"]}} + authored = AgenticComponentActivity( + name="Authored", + task_key="authored", + notifications=collapsed, + resources=RESOURCES, + task={**TASK, "email_notifications": {"on_success": ["owner@example.com"]}}, + ) + plain = AgenticComponentActivity(name="Plain", task_key="plain", notifications=collapsed, task=TASK) + + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[authored, plain])), tmp_path) + + job_resource = yaml.safe_load((tmp_path / "resources" / "orders.yml").read_text(encoding="utf-8")) + tasks = {task["task_key"]: task for task in job_resource["resources"]["jobs"]["orders"]["tasks"]} + assert tasks["authored"]["email_notifications"] == {"on_success": ["owner@example.com"]} + assert tasks["plain"]["email_notifications"] == {"on_failure": ["flowx@example.com"]} + + +def test_agentic_component_with_its_own_compute_never_gets_a_second_binding(tmp_path): + serverless = AgenticComponentActivity( + name="Serverless", + task_key="serverless", + environments=[{**SERVERLESS_ENVIRONMENT, "environment_key": "default"}], + task={"notebook_task": {"notebook_path": "/Workspace/Shared/etl/orders"}, "environment_key": "default"}, + ) + unbound = AgenticComponentActivity( + name="Unbound", + task_key="unbound", + task={"notebook_task": {"notebook_path": "/Workspace/Shared/etl/customers"}}, + ) + + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[serverless, unbound])), tmp_path) + + job_resource = yaml.safe_load((tmp_path / "resources" / "orders.yml").read_text(encoding="utf-8")) + tasks = {task["task_key"]: task for task in job_resource["resources"]["jobs"]["orders"]["tasks"]} + assert tasks["serverless"] == { + "task_key": "serverless", + "notebook_task": {"notebook_path": "/Workspace/Shared/etl/orders"}, + "environment_key": "default", + } + assert tasks["unbound"]["job_cluster_key"] == "default_cluster" + assert job_resource["resources"]["jobs"]["orders"]["environments"] == [ + {**SERVERLESS_ENVIRONMENT, "environment_key": "default"} + ] + assert check_bundle_dir(tmp_path).ok + + +def test_agentic_component_base_parameters_that_look_dynamic_are_written_as_given(tmp_path): + activity = AgenticComponentActivity( + name="Report", + task_key="report", + files=[{"path": "notebooks/report.py", "content": "print('report')\n"}], + task={ + "notebook_task": { + "notebook_path": "../src/notebooks/report.py", + "base_parameters": {"cutoff": "datetime.now(timezone.utc) - timedelta(days=7)"}, + } + }, + ) + + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[activity])), tmp_path) + + job_resource = yaml.safe_load((tmp_path / "resources" / "orders.yml").read_text(encoding="utf-8")) + assert job_resource["resources"]["jobs"]["orders"]["tasks"] == [ + { + "task_key": "report", + "notebook_task": { + "notebook_path": "../src/notebooks/report.py", + "base_parameters": {"cutoff": "datetime.now(timezone.utc) - timedelta(days=7)"}, + }, + } + ] + + +def test_agentic_component_environments_round_trip_through_serialized_ir(): + activity = AgenticComponentActivity( + name="Run", + task_key="run", + environments=[SERVERLESS_ENVIRONMENT], + task={"spark_python_task": {"python_file": "../src/jobs/run.py"}, "environment_key": "serverless"}, + ) + + serialized = json.loads(json.dumps(activity_to_dict(activity))) + restored_pipeline, _ = pipeline_dict_to_ir({"name": "orders", "tasks": [serialized]}) + + assert serialized["environments"] == [SERVERLESS_ENVIRONMENT] + assert restored_pipeline.tasks[0].environments == [SERVERLESS_ENVIRONMENT] + assert activity_to_dict(restored_pipeline.tasks[0]) == serialized + + +def _serverless_component(payload_kind: str, task_key: str = "run") -> AgenticComponentActivity: + if payload_kind == "spark_python_task": + files = [{"path": "jobs/run.py", "content": "print('run')\n"}] + payload = {"spark_python_task": {"python_file": "../src/jobs/run.py"}} + environment = SERVERLESS_ENVIRONMENT + else: + files = [{"path": "libraries/orders.whl", "binary_content": "UEsDBAoAAAAA"}] + payload = {"python_wheel_task": {"package_name": "orders", "entry_point": "main"}} + environment = WHEEL_ENVIRONMENT + return AgenticComponentActivity( + name=task_key.title(), + task_key=task_key, + files=files, + environments=[environment], + task={**payload, "environment_key": "serverless"}, + ) + + +@pytest.mark.parametrize("payload_kind", ["spark_python_task", "python_wheel_task"]) +def test_serverless_agentic_component_gets_its_environment_in_the_job(tmp_path, payload_kind): + component = _serverless_component(payload_kind) + serialized_pipeline = json.loads(json.dumps(pipeline_to_dict(Pipeline(name="orders", tasks=[component])))) + restored_pipeline, _ = pipeline_dict_to_ir(serialized_pipeline) + + write_bundle(prepare_workflow(restored_pipeline), tmp_path) + + job = yaml.safe_load((tmp_path / "resources" / "orders.yml").read_text(encoding="utf-8"))["resources"]["jobs"][ + "orders" + ] + assert job["environments"] == component.environments + assert job["tasks"] == [{"task_key": "run", **component.task}] + assert "job_clusters" not in job + if payload_kind == "python_wheel_task": + assert job["environments"][0]["spec"]["dependencies"] == ["../src/libraries/orders.whl"] + assert (tmp_path / "src" / "libraries" / "orders.whl").exists() + assert check_bundle_dir(tmp_path).ok + + +def _jobs_by_task_key(bundle_dir) -> dict[str, dict]: + """Map every emitted task key (including ForEach bodies) to the job that holds it.""" + jobs_by_task_key: dict[str, dict] = {} + for resource_path in sorted((bundle_dir / "resources").glob("*.yml")): + jobs = (yaml.safe_load(resource_path.read_text(encoding="utf-8"))["resources"]).get("jobs") or {} + for job in jobs.values(): + pending = list(job["tasks"]) + while pending: + task = pending.pop() + jobs_by_task_key[task["task_key"]] = job + body = (task.get("for_each_task") or {}).get("task") + if body: + pending.append(body) + return jobs_by_task_key + + +@pytest.mark.parametrize( + "wrap", [_inside_if_condition, _inside_switch, _inside_for_each, _inside_for_each_with_siblings] +) +def test_serverless_agentic_component_environment_lands_in_the_job_that_runs_it(tmp_path, wrap): + write_bundle( + prepare_workflow(Pipeline(name="orders", tasks=[wrap(_serverless_component("spark_python_task"))])), + tmp_path, + ) + + jobs_by_task_key = _jobs_by_task_key(tmp_path) + assert jobs_by_task_key["run"]["environments"] == [SERVERLESS_ENVIRONMENT] + assert check_bundle_dir(tmp_path).ok + + +def test_agentic_component_environment_key_that_no_component_declares_fails(tmp_path): + activity = AgenticComponentActivity( + name="Run", + task_key="run", + files=[{"path": "jobs/run.py", "content": "print('run')\n"}], + task={"spark_python_task": {"python_file": "../src/jobs/run.py"}, "environment_key": "missing"}, + ) + + with pytest.raises(ValueError, match="environment_key 'missing'"): + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[activity])), tmp_path) + assert not (tmp_path / "databricks.yml").exists() + + +def test_agentic_components_declaring_an_identical_environment_write_it_once(tmp_path): + first = _serverless_component("spark_python_task", task_key="first") + second = _serverless_component("spark_python_task", task_key="second") + + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[first, second])), tmp_path) + + job = yaml.safe_load((tmp_path / "resources" / "orders.yml").read_text(encoding="utf-8"))["resources"]["jobs"][ + "orders" + ] + assert job["environments"] == [SERVERLESS_ENVIRONMENT] + assert check_bundle_dir(tmp_path).ok + + +def test_agentic_components_declaring_one_environment_key_with_different_specs_fail(tmp_path): + first = _serverless_component("spark_python_task", task_key="first") + second = _serverless_component("spark_python_task", task_key="second") + second.environments = [{"environment_key": "serverless", "spec": {"environment_version": "3"}}] + + with pytest.raises(ValueError, match="declared more than once with different specs"): + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[first, second])), tmp_path) + assert not (tmp_path / "databricks.yml").exists() + + +@pytest.mark.parametrize( + "environment", + [ + "serverless", + {"environment_key": "serverless"}, + {"environment_key": "serverless", "spec": "2"}, + {"environment_key": "", "spec": {"environment_version": "2"}}, + {"environment_key": 5, "spec": {"environment_version": "2"}}, + {**SERVERLESS_ENVIRONMENT, "comment": "extra"}, + ], +) +def test_agentic_component_rejects_malformed_environments(environment): + activity = AgenticComponentActivity( + name="Bad", + task_key="bad", + environments=[environment], + task={"spark_python_task": {"python_file": "../src/jobs/run.py"}, "environment_key": "serverless"}, + ) + + with pytest.raises(ValueError, match="Agentic component 'bad' environment"): + prepare_workflow(Pipeline(name="orders", tasks=[activity])) + + +@pytest.mark.parametrize( + "payload", + [ + {"spark_python_task": {"python_file": "../src/jobs/run.py"}}, + {"python_wheel_task": {"package_name": "orders", "entry_point": "main"}}, + {"spark_jar_task": {"main_class_name": "com.example.Main"}}, + {"spark_submit_task": {"parameters": ["--class", "com.example.Main", "dbfs:/jars/etl.jar"]}}, + {"dbt_task": {"commands": ["dbt deps", "dbt run"], "warehouse_id": "abc123"}}, + { + "for_each_task": { + "inputs": "[1, 2]", + "task": {"task_key": "run_item", "spark_python_task": {"python_file": "../src/jobs/run.py"}}, + } + }, + ], +) +def test_agentic_component_python_or_jar_task_without_compute_fails(payload): + """flowx binds a cluster only for notebooks, so these would package with no compute and fail at deploy.""" + activity = AgenticComponentActivity(name="Run", task_key="run", task=payload) + + with pytest.raises(ValueError, match="'run' .* must name its compute"): + prepare_workflow(Pipeline(name="orders", tasks=[activity])) + + +@pytest.mark.parametrize( + "compute", [{"environment_key": "serverless"}, {"existing_cluster_id": "0101-123456-abcdefgh"}] +) +def test_agentic_component_spark_submit_task_without_a_new_cluster_fails(compute): + """Databricks runs spark-submit only on a new cluster, so serverless or an existing cluster fails at deploy.""" + activity = AgenticComponentActivity( + name="Run", + task_key="run", + environments=[SERVERLESS_ENVIRONMENT], + task={"spark_submit_task": {"parameters": ["--class", "com.example.Main", "dbfs:/jars/etl.jar"]}, **compute}, + ) + + with pytest.raises( + ValueError, match="'run' spark_submit_task must name its compute with job_cluster_key, new_cluster" + ): + prepare_workflow(Pipeline(name="orders", tasks=[activity])) + + +@pytest.mark.parametrize( + "task", + [ + {"spark_python_task": {"python_file": "../src/jobs/run.py"}, "job_cluster_key": "etl_cluster"}, + {"notebook_task": {"notebook_path": "/Workspace/Shared/etl/orders"}, "job_cluster_key": "etl_cluster"}, + { + "for_each_task": { + "inputs": "[1, 2]", + "task": { + "task_key": "run_item", + "spark_python_task": {"python_file": "../src/jobs/run.py"}, + "job_cluster_key": "etl_cluster", + }, + } + }, + ], +) +def test_agentic_component_job_cluster_key_the_bundle_does_not_define_fails(task): + """A component cannot declare a job cluster, so any other key would package bound to nothing.""" + activity = AgenticComponentActivity(name="Run", task_key="run", task=task) + + with pytest.raises(ValueError, match="'run' job_cluster_key 'etl_cluster' must be one of the bundle's job"): + prepare_workflow(Pipeline(name="orders", tasks=[activity])) + + +@pytest.mark.parametrize("cluster_key", ["default_cluster", "single_node_cluster", "multi_node_cluster"]) +@pytest.mark.parametrize("wrap", [lambda activity: activity, _inside_for_each, _inside_for_each_with_siblings]) +def test_agentic_component_bound_to_a_bundle_job_cluster_gets_that_cluster_in_its_job(tmp_path, cluster_key, wrap): + component = AgenticComponentActivity( + name="Run", + task_key="run", + files=[{"path": "jobs/run.py", "content": "print('run')\n"}], + task={"spark_python_task": {"python_file": "../src/jobs/run.py"}, "job_cluster_key": cluster_key}, + ) + + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[wrap(component)])), tmp_path) + + job = _jobs_by_task_key(tmp_path)["run"] + assert [cluster["job_cluster_key"] for cluster in job["job_clusters"]] == [cluster_key] + assert check_bundle_dir(tmp_path).ok + + +def test_agentic_component_file_without_content_or_binary_content_fails(): + activity = AgenticComponentActivity( + name="Run", + task_key="run", + files=[{"path": "notebooks/run.py"}], + task={"notebook_task": {"notebook_path": "../src/notebooks/run.py"}}, + ) + + with pytest.raises(ValueError, match="'run' file 'notebooks/run.py' must set content or binary_content"): + prepare_workflow(Pipeline(name="orders", tasks=[activity])) + + +@pytest.mark.parametrize("content", ["print('no trailing newline')", ""]) +def test_agentic_component_text_file_is_written_byte_for_byte(tmp_path, content): + activity = AgenticComponentActivity( + name="Run", + task_key="run", + files=[{"path": "notebooks/run.py", "content": content}], + task={"notebook_task": {"notebook_path": "../src/notebooks/run.py"}}, + ) + + write_bundle(prepare_workflow(Pipeline(name="orders", tasks=[activity])), tmp_path) + + assert (tmp_path / "src" / "notebooks" / "run.py").read_bytes() == content.encode("utf-8") diff --git a/tests/unit/test_bundle_invariants.py b/tests/unit/test_bundle_invariants.py index 8b52e55c..5bb4328e 100644 --- a/tests/unit/test_bundle_invariants.py +++ b/tests/unit/test_bundle_invariants.py @@ -130,3 +130,82 @@ def test_bundle_job_reference_to_unknown_resource_is_flagged(tmp_path): assert "parent.yml" in finding.location assert "call_missing" in finding.location assert "dangling_run_job_reference" in format_result(result) + + +def test_bundle_pipeline_reference_can_target_pipeline_in_another_resource_file(tmp_path): + resources = tmp_path / "resources" + resources.mkdir() + (resources / "job.yml").write_text( + "resources:\n" + " jobs:\n" + " parent:\n" + " tasks:\n" + " - task_key: run_pipeline\n" + " pipeline_task:\n" + " pipeline_id: ${resources.pipelines.ingestion.id}\n", + encoding="utf-8", + ) + (resources / "pipeline.yml").write_text( + "resources:\n pipelines:\n ingestion:\n name: ingestion\n", + encoding="utf-8", + ) + + result = check_bundle_dir(tmp_path) + + assert "dangling_pipeline_reference" not in _codes(result.findings) + + +def test_bundle_pipeline_reference_to_unknown_resource_is_flagged(tmp_path): + resources = tmp_path / "resources" + resources.mkdir() + (resources / "job.yml").write_text( + "resources:\n" + " jobs:\n" + " parent:\n" + " tasks:\n" + " - task_key: run_pipeline\n" + " pipeline_task:\n" + " pipeline_id: ${resources.pipelines.missing.id}\n", + encoding="utf-8", + ) + + result = check_bundle_dir(tmp_path) + finding = next(finding for finding in result.findings if finding.code == "dangling_pipeline_reference") + + assert finding.severity == "violation" + assert "job.yml" in finding.location + assert "run_pipeline" in finding.location + assert "dangling_pipeline_reference" in format_result(result) + + +def test_bundle_environment_reference_must_be_declared_in_the_same_job(tmp_path): + resources = tmp_path / "resources" + resources.mkdir() + (resources / "job.yml").write_text( + "resources:\n" + " jobs:\n" + " declared:\n" + " environments:\n" + " - environment_key: serverless\n" + " spec:\n" + " environment_version: '2'\n" + " tasks:\n" + " - task_key: run_declared\n" + " environment_key: serverless\n" + " spark_python_task:\n" + " python_file: ../src/run.py\n" + " undeclared:\n" + " tasks:\n" + " - task_key: run_undeclared\n" + " environment_key: serverless\n" + " spark_python_task:\n" + " python_file: ../src/run.py\n", + encoding="utf-8", + ) + + result = check_bundle_dir(tmp_path) + findings = [finding for finding in result.findings if finding.code == "dangling_environment_reference"] + + assert len(findings) == 1 + assert findings[0].severity == "violation" + assert "run_undeclared" in findings[0].location