diff --git a/README.md b/README.md index c930128..bf638d2 100644 --- a/README.md +++ b/README.md @@ -511,6 +511,34 @@ def greeting_pipeline(who: In[str], cfg) -> Out[str]: Task IDs default from the left-hand variable name at the call site, converted to title case. If there is no simple left-hand variable, or if you want a stable explicit label, call `.named("Task Id")` before invoking the task. Use `.bind(...)` to pre-fill task arguments and `.with_annotations({...})` to add per-task annotations. +##### Root pipeline annotations + +`@pipeline(annotations={...})` writes the compiled pipeline's root `metadata.annotations` block. A caller that compiles programmatically can supply the block instead — typically from its own per-environment config file, so the values do not have to be hard-coded in source: + +```python +from tangle_cli.pipelines import compile_pipeline_file + +compile_pipeline_file( + "pipeline.py", + "pipeline.yaml", + pipeline_annotations={"environment": "staging", "owner": "search-platform"}, +) +``` + +The same keyword exists on `tangle_cli.pipeline_compiler.compile_pipeline` and on `PipelineCompiler.compile_file`. There is no CLI flag: the source route already exists, and what the keyword adds is a programmatic/config route for the part of the block that varies by environment. + +Semantics: + +- **Per-key merge, caller wins.** `@pipeline(annotations={"author": "a", "version": "1.0"})` compiled with `pipeline_annotations={"version": "2.0", "environment": "staging"}` emits all three keys, with `version: "2.0"`. Source keys the caller does not mention are preserved, so invariants stay in source and only the varying subset is passed in. +- **Omitted or `{}` is a no-op**, byte for byte — an empty mapping is not a destructive clear of the source block. +- **Root only.** `subpipeline` children never inherit it, so child subgraph sidecar names, bytes, and component digests are unaffected. A child that wants annotations declares its own. +- **Descriptive only.** Root metadata is not read by the orchestrator, so it cannot influence placement, routing, scheduling, or run identity. Use pipeline-run annotations for anything execution-bearing. +- **`str -> str`, validated up front.** A non-mapping argument, a non-string key or value, an empty key, a `system/`-prefixed key (reserved by Tangle), or a template delimiter (`{{`, `{%`, `{#`) in a key or value raises `InvalidPipelineAnnotationsError` (a `CompileError`) before anything is imported or written. Annotations usually come from an untrusted config file, so the diagnostics name the key and the type and never echo a value. Values are baked into the compiled YAML and the stored pipeline definition: labels only, never secrets. + +The rules live in one place, `tangle_cli.schema_validation`: `check_annotations(mapping, policy=..., error_cls=...)` applied under a named `AnnotationPolicy`. `CALLER_ANNOTATION_POLICY` is the strict input policy described above; `DOCUMENT_ANNOTATION_POLICY` is the lenient policy every pipeline document is validated against (scalar-or-null values, no key rules), matching the schema and hand-authored YAML. Only the caller-supplied input surface is strict: existing documents are accepted exactly as before, and the document check still runs on the merged result, so annotations reaching the output by any route are validated. + +A distribution that reads these annotations from its own config file should call `check_annotations(mapping, policy=CALLER_ANNOTATION_POLICY, error_cls=...)` at config-parse time — passing its own error type and adding the config path and key to the message — so one user mistake produces one diagnostic instead of two competing ones. The compiler's own call is then the backstop for anything arriving by another route. + ##### Conditional task execution Pipeline inputs used as conditions are ordinary `In[str]` values; there is no special conditional input annotation. Pass the value through the reserved task-call metadata keyword `is_enabled=`: diff --git a/packages/tangle-cli/src/tangle_cli/__init__.py b/packages/tangle-cli/src/tangle_cli/__init__.py index 926f528..4eb9746 100644 --- a/packages/tangle-cli/src/tangle_cli/__init__.py +++ b/packages/tangle-cli/src/tangle_cli/__init__.py @@ -14,6 +14,6 @@ try: __version__ = metadata_version("tangle-cli") except PackageNotFoundError: - __version__ = "0.1.18" + __version__ = "0.1.19" __all__ = ["TangleDynamicDiscoveryClient", "__version__"] diff --git a/packages/tangle-cli/src/tangle_cli/pipeline_compiler.py b/packages/tangle-cli/src/tangle_cli/pipeline_compiler.py index b7365d4..8b7040e 100644 --- a/packages/tangle-cli/src/tangle_cli/pipeline_compiler.py +++ b/packages/tangle-cli/src/tangle_cli/pipeline_compiler.py @@ -65,14 +65,19 @@ overrides_fingerprint, ) from .python_pipeline.emit import _TASK_URL_PLACEHOLDER, emit_pipeline -from .python_pipeline.errors import CompileError +from .python_pipeline.errors import CompileError, InvalidPipelineAnnotationsError from .python_pipeline.pipeline import PipelineFn from .python_pipeline.ref import CallableRef from .python_pipeline.registered import _REGISTERED_URL_PLACEHOLDER from .python_pipeline.subpipeline import _SUBPIPELINE_URL_PLACEHOLDER, SubpipelineRef from .python_pipeline.trace import trace_pipeline from .python_pipeline.types import In -from .schema_validation import SchemaValidationError, validate_dehydrated_pipeline +from .schema_validation import ( + CALLER_ANNOTATION_POLICY, + SchemaValidationError, + check_annotations, + validate_dehydrated_pipeline, +) from .utils import dump_yaml @@ -141,6 +146,7 @@ def compile_pipeline( pipeline_name: str | None = None, emit_components_sidecar: bool = True, image_overrides: Mapping[str, str] | None = None, + pipeline_annotations: Mapping[str, str] | None = None, ) -> CompileResult: """Compile ``script`` to a single pipeline YAML at ``output``. @@ -174,6 +180,18 @@ def compile_pipeline( image_overrides: Compile-time image-id overrides from ``--image ID=REF``. They apply only to ``@task(image_id=...)`` refs that do not also set an explicit ``image=``. + pipeline_annotations: Caller-supplied ROOT ``metadata.annotations`` + (``str -> str``), typically read from a downstream config file so + the block can vary per environment. Merged PER KEY over the root + ``@pipeline(annotations=...)`` block, caller winning on collision; + ``None`` / ``{}`` is a no-op that leaves the compiled bytes + identical. ROOT ONLY — ``subpipeline`` children do not inherit it, + so their sidecar names, bytes, and digests are unaffected. Root + metadata is descriptive: it is not read by the orchestrator and + cannot influence placement, routing, scheduling, or identity. See + :data:`tangle_cli.schema_validation.CALLER_ANNOTATION_POLICY` + for the accepted shape; a malformed mapping raises + :class:`~tangle_cli.python_pipeline.errors.InvalidPipelineAnnotationsError`. Returns: A :class:`CompileResult`. ``components_path`` is the sidecar path @@ -186,6 +204,15 @@ def compile_pipeline( """ overrides = dict(overrides or {}) image_overrides = dict(image_overrides or {}) + # Validated up front, by the shared validation layer, so a hostile or + # malformed annotation mapping fails before any module is imported or any + # file is written. The compiler owns no annotation rules of its own; it + # only chooses the caller-facing error type. + root_annotations = check_annotations( + pipeline_annotations, + policy=CALLER_ANNOTATION_POLICY, + error_cls=InvalidPipelineAnnotationsError, + ) # 1. Validate the script path. script_path = Path(script).resolve() @@ -241,6 +268,7 @@ def compile_pipeline( emit_components_sidecar=emit_components_sidecar, source_dirs=purge_dirs, image_overrides=image_overrides, + pipeline_annotations=root_annotations, ) # 5. Compile the root (and, recursively, all children) into in-memory @@ -445,6 +473,12 @@ def _compile_pipeline_fn( # output guard must skip for THIS artifact's body. with _temp_sys_path(base_dir): builder = trace_pipeline(pipeline_fn, cfg=cfg, inputs={}) + # Caller-supplied root annotations: ROOT ONLY (a child keeps exactly what + # its own ``@pipeline`` declared), merged PER KEY with the caller winning + # on collision, before emit so they pass the same emit-time guards as + # authored ones. ``update`` keeps the emitted order source-first. + if is_root and ctx.pipeline_annotations: + builder.annotations.update(ctx.pipeline_annotations) body_dict, exempt_paths = emit_pipeline(builder) # 3. The canonical compile key for dedup / cycle detection. A child is @@ -2785,6 +2819,7 @@ def compile_file( pipeline_name: str | None = None, emit_components_sidecar: bool = True, image_overrides: Mapping[str, str] | None = None, + pipeline_annotations: Mapping[str, str] | None = None, ) -> CompileResult: """Compile ``script`` to a single dehydrated pipeline YAML at ``output``. @@ -2804,6 +2839,7 @@ def compile_file( pipeline_name=pipeline_name, emit_components_sidecar=emit_components_sidecar, image_overrides=image_overrides, + pipeline_annotations=pipeline_annotations, ) self.log.info(f"wrote {result.pipeline_path}") if result.components_path is not None: diff --git a/packages/tangle-cli/src/tangle_cli/pipelines.py b/packages/tangle-cli/src/tangle_cli/pipelines.py index d21f67d..5339dc1 100644 --- a/packages/tangle-cli/src/tangle_cli/pipelines.py +++ b/packages/tangle-cli/src/tangle_cli/pipelines.py @@ -259,6 +259,7 @@ def compile_pipeline_file( pipeline_name: str | None = None, emit_components_sidecar: bool = True, image_overrides: Mapping[str, str] | None = None, + pipeline_annotations: Mapping[str, str] | None = None, logger: Any | None = None, ) -> CompileResult: """Compile a Python-authored pipeline to a dehydrated YAML bundle. @@ -271,6 +272,12 @@ def compile_pipeline_file( The :class:`~tangle_cli.pipeline_compiler.CompileResult` is returned as-is — unlike hydrate, the compiler already exposes its public result type, so there is nothing to repackage. + + ``pipeline_annotations`` supplies ROOT ``metadata.annotations`` (``str -> + str``) merged per key over the root ``@pipeline(annotations=...)`` block, + caller winning on collision; it applies to the root only and is a no-op + when omitted or empty. See + :func:`~tangle_cli.pipeline_compiler.compile_pipeline`. """ from .pipeline_compiler import PipelineCompiler @@ -286,6 +293,10 @@ def compile_pipeline_file( pipeline_name=pipeline_name, emit_components_sidecar=emit_components_sidecar, image_overrides=dict(image_overrides) if image_overrides else None, + # Passed through unconverted: the compiler validates the shape and + # takes its own copy, so a malformed mapping fails with the + # compiler's value-free CompileError instead of a bare TypeError. + pipeline_annotations=pipeline_annotations, ) except (CompileError, SchemaValidationError) as exc: raise PipelineValidationError(str(exc)) from exc diff --git a/packages/tangle-cli/src/tangle_cli/python_pipeline/compiler_context.py b/packages/tangle-cli/src/tangle_cli/python_pipeline/compiler_context.py index a54d270..118620f 100644 --- a/packages/tangle-cli/src/tangle_cli/python_pipeline/compiler_context.py +++ b/packages/tangle-cli/src/tangle_cli/python_pipeline/compiler_context.py @@ -253,6 +253,14 @@ class CompileContext: root_overrides: dict[str, str] = field(default_factory=dict) emit_components_sidecar: bool = True image_overrides: dict[str, str] = field(default_factory=dict) + # Caller-supplied ROOT ``metadata.annotations`` (already validated), + # merged PER KEY over the root ``@pipeline(annotations=...)`` block just + # before emit. ROOT ONLY: children never inherit it, so a child sidecar's + # bytes — and therefore its component digest — are untouched by it. It is + # deliberately absent from :class:`PipelineCompileKey` / + # ``overrides_fingerprint`` (as ``image_overrides`` is), so sidecar + # filenames stay identity-derived rather than content-derived. + pipeline_annotations: dict[str, str] = field(default_factory=dict) max_depth: int = 32 # Compiled CHILD artifacts keyed by compile key (Decision M dedup). # The root is NOT stored here; it is returned directly. diff --git a/packages/tangle-cli/src/tangle_cli/python_pipeline/errors.py b/packages/tangle-cli/src/tangle_cli/python_pipeline/errors.py index 19b684f..66f3bf6 100644 --- a/packages/tangle-cli/src/tangle_cli/python_pipeline/errors.py +++ b/packages/tangle-cli/src/tangle_cli/python_pipeline/errors.py @@ -27,3 +27,14 @@ class AmbiguousTaskIdError(CompileError): class InvalidArgumentTypeError(CompileError): """Raised on an argument value with no supported emit dispatch.""" + + +class InvalidPipelineAnnotationsError(CompileError): + """Raised on a malformed caller-supplied ``pipeline_annotations`` mapping. + + A dedicated type because these annotations usually originate in a + downstream CONFIG file: a caller that reads such a file can catch this + precisely and re-raise with the config path and key attached, without + broadly catching :class:`CompileError` and swallowing unrelated compile + failures. Messages never echo an annotation value. + """ diff --git a/packages/tangle-cli/src/tangle_cli/schema_validation.py b/packages/tangle-cli/src/tangle_cli/schema_validation.py index fffc3c2..ded305e 100644 --- a/packages/tangle-cli/src/tangle_cli/schema_validation.py +++ b/packages/tangle-cli/src/tangle_cli/schema_validation.py @@ -25,6 +25,12 @@ ``taskOutput.taskId``, undeclared ``graphInput.inputName``, ``outputValues`` ↔ ``outputs`` correspondence, scalar metadata annotations, pure componentRefs) and the no-template-delimiter scan. +* :func:`check_annotations` plus :data:`DOCUMENT_ANNOTATION_POLICY` / + :data:`CALLER_ANNOTATION_POLICY` — the ONE place that says what a + ``metadata.annotations`` mapping may contain. Both the lenient check + applied to every pipeline DOCUMENT and the strict check applied to a + caller/config-supplied mapping are the same function under different + policies, so the two can never drift apart. Everything here is standalone — it never changes ``PipelineHydrator`` behavior. ``compile_pipeline`` uses :func:`validate_dehydrated_pipeline` @@ -34,6 +40,7 @@ import json from collections.abc import Collection, Iterator, Mapping +from dataclasses import dataclass from functools import lru_cache from pathlib import Path from typing import Any @@ -195,6 +202,152 @@ def assert_no_template_delimiters( ) +# --------------------------------------------------------------------------- +# Annotation policy: the one definition of what an annotations mapping may +# contain, applied at two documented strictness levels. + +#: Annotation-key prefix Tangle reserves for its own annotations. +RESERVED_ANNOTATION_KEY_PREFIX = "system/" + + +@dataclass(frozen=True) +class AnnotationPolicy: + """Which rules :func:`check_annotations` applies, and how it names the + mapping in diagnostics (``label``). + + Every flag defaults to off, so a policy opts IN to strictness. With + ``require_string_values`` off, values fall back to the legacy + scalar-or-null rule rather than being unchecked. + """ + + label: str + require_string_values: bool = False + require_string_keys: bool = False + require_non_empty_keys: bool = False + reject_reserved_key_prefix: bool = False + reject_template_delimiters: bool = False + reject_non_mapping: bool = False + + +#: Applied to every pipeline document validated by +#: :func:`validate_dehydrated_pipeline`. Deliberately lenient: hand-authored +#: and legacy YAML reaches this path, and any flag enabled here would reject +#: documents that compile and run today. +DOCUMENT_ANNOTATION_POLICY = AnnotationPolicy(label="metadata.annotations") + +#: Applied to a caller/config-supplied root annotations mapping (the +#: ``pipeline_annotations`` compile argument, and any downstream config reader +#: that wants to refuse a bad mapping early with its own file/key provenance: +#: ``check_annotations(mapping, policy=CALLER_ANNOTATION_POLICY, +#: error_cls=...)``). Strict ``str -> str`` because that input surface is new, +#: has no legacy, and is usually an untrusted config file — hence also +#: :func:`check_annotations`'s rule that a diagnostic names the key and the +#: type or delimiter but never echoes a value. ``None`` / ``{}`` means nothing +#: supplied and is a no-op, not a clear. ``reject_template_delimiters`` adds no +#: rule the output lacks — :func:`assert_no_template_delimiters` scans the +#: compiled output regardless; it moves the failure earlier, onto a message +#: that names the offending annotation key. +CALLER_ANNOTATION_POLICY = AnnotationPolicy( + label="pipeline_annotations", + require_string_values=True, + require_string_keys=True, + require_non_empty_keys=True, + reject_reserved_key_prefix=True, + reject_template_delimiters=True, + reject_non_mapping=True, +) + + +def check_annotations( + annotations: Any, + *, + policy: AnnotationPolicy, + error_cls: type[Exception] = SchemaValidationError, +) -> dict[Any, Any]: + """Check an annotations mapping against ``policy`` and copy it. + + Annotations are frequently UNTRUSTED input (a config file, a checked-in + YAML document), so every diagnostic names the offending KEY and the + offending TYPE or delimiter and NEVER echoes a value. + + Args: + annotations: The mapping to check. ``None`` means "nothing supplied" + and is always accepted. + policy: Which rules apply — :data:`DOCUMENT_ANNOTATION_POLICY` or + :data:`CALLER_ANNOTATION_POLICY`. + error_cls: Exception type to raise, so a caller-input path can raise + its own precise type (``InvalidPipelineAnnotationsError``) while + the document path keeps raising + :class:`SchemaValidationError`. This module deliberately does + not import the authoring-layer error hierarchy. + + Returns: + A plain ``dict`` copy in the caller's key order (empty for ``None``), + so later mutation of the supplied mapping cannot reach the compile. + No value is coerced or normalized — this function only accepts or + rejects. + + Raises: + error_cls: on the first violation found. + """ + if annotations is None: + return {} + if not isinstance(annotations, Mapping): + if not policy.reject_non_mapping: + return {} + raise error_cls( + f"{policy.label} must be a mapping of string keys to string " + f"values; got {type(annotations).__name__}." + ) + + checked: dict[Any, Any] = {} + for key, value in annotations.items(): + if policy.require_string_keys and not isinstance(key, str): + raise error_cls( + f"{policy.label} keys must be strings; got a key of type " + f"{type(key).__name__}." + ) + if policy.require_non_empty_keys and key == "": + raise error_cls(f"{policy.label} keys must not be empty.") + if policy.require_string_values: + # Report the type only — a rejected value may be untrusted. + if not isinstance(value, str): + raise error_cls( + f"{policy.label} value for key {key!r} must be a string; " + f"got {type(value).__name__}." + ) + # bool is a subclass of int, so it is covered by int. + elif not (value is None or isinstance(value, (str, int, float))): + raise error_cls( + f"{policy.label}[{key!r}] must be a scalar (str/number/bool) " + f"or null, got {type(value).__name__!r}." + ) + if policy.reject_reserved_key_prefix and isinstance(key, str): + if key.startswith(RESERVED_ANNOTATION_KEY_PREFIX): + raise error_cls( + f"{policy.label} key {key!r} uses the reserved " + f"{RESERVED_ANNOTATION_KEY_PREFIX!r} prefix, which Tangle " + "keeps for its own annotations. Choose a different key." + ) + if policy.reject_template_delimiters: + # Scanned with the SAME tokens the compiled-output guard uses, so + # an input check and the output contract can never disagree. + for where, text in ((f"key {key!r}", key), (f"value for key {key!r}", value)): + if not isinstance(text, str): + continue + delim = next((d for _p, d in iter_template_delimiters(text)), None) + if delim is None: + continue + raise error_cls( + f"{policy.label} {where} contains the template delimiter " + f"{delim!r}. Annotations are literal metadata written " + "verbatim into the compiled pipeline, which must be fully " + "rendered — resolve the template before passing the value." + ) + checked[key] = value + return checked + + # --------------------------------------------------------------------------- # Phase 5: dehydrated-pipeline shape detection + semantic validation. # @@ -411,15 +564,12 @@ def _validate_semantics(data: Mapping[str, Any]) -> None: metadata = data.get("metadata") if isinstance(metadata, Mapping): annotations = metadata.get("annotations") + # A non-mapping ``annotations`` is left to JSON-Schema, exactly as + # before. Everything a DOCUMENT's annotations may contain is decided + # by the shared policy below, so this backstop and the strict + # caller-input check can never disagree about a shared rule. if isinstance(annotations, Mapping): - for key, value in annotations.items(): - # bool is a subclass of int, so it is covered by int. - if not (value is None or isinstance(value, (str, int, float))): - raise SchemaValidationError( - f"metadata.annotations[{key!r}] must be a scalar " - f"(str/number/bool) or null, got " - f"{type(value).__name__!r}." - ) + check_annotations(annotations, policy=DOCUMENT_ANNOTATION_POLICY) def validate_dehydrated_pipeline( diff --git a/pyproject.toml b/pyproject.toml index 66fe108..41e2f02 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "tangle-cli" -version = "0.1.18" +version = "0.1.19" description = "CLI for Tangle, the open-source ML pipeline orchestration platform" readme = "README.md" authors = [ diff --git a/tests/test_packaging.py b/tests/test_packaging.py index 6534c7e..9c6965e 100644 --- a/tests/test_packaging.py +++ b/tests/test_packaging.py @@ -183,7 +183,7 @@ def test_tangle_cli_wheel_supports_expert_no_deps_import_path_without_tangle_api requires_dist = [line for line in metadata.splitlines() if line.startswith("Requires-Dist: ")] assert not any(name.startswith("tangle_api/") for name in names) assert "tangle_cli/openapi/openapi.json" not in names - assert "Version: 0.1.18" in metadata + assert "Version: 0.1.19" in metadata assert "Requires-Dist: tangle-api==0.1.1" in requires_dist assert not any("extra == 'native'" in line for line in requires_dist) assert "Provides-Extra: native" in metadata diff --git a/tests/test_pipeline_annotations.py b/tests/test_pipeline_annotations.py new file mode 100644 index 0000000..baa6f88 --- /dev/null +++ b/tests/test_pipeline_annotations.py @@ -0,0 +1,586 @@ +"""Caller-supplied ROOT ``metadata.annotations`` (``pipeline_annotations``). + +``compile_pipeline(..., pipeline_annotations={...})`` writes the compiled +pipeline's ROOT ``metadata.annotations`` block, so a downstream caller (in +practice a per-environment config file) can set descriptive metadata that the +authoring source does not hard-code. The contract these tests pin: + +* **Per-key merge, caller wins.** Source ``@pipeline(annotations=...)`` keys the + caller does not mention survive; a collision resolves to the caller's value. +* **Absent / empty is a no-op**, proven by BYTE identity of the whole bundle — + ``{}`` is not a destructive clear of the source block. +* **Root only.** A ``subpipeline`` child inherits nothing, so its sidecar + filename (``-.yaml``, hashed over compile IDENTITY) and its + bytes — and therefore its component digest — are untouched. +* **Hostile input is refused at parse time**, before any module is imported or + any file is written, with a diagnostic that names the key and the type but + NEVER echoes a value. +* **The value survives hydration**, which rewrites componentRefs and must not + disturb root metadata. + +One more thing is pinned here: annotation rules live in EXACTLY ONE place, +``schema_validation.check_annotations``, applied under two documented +policies — the same entry point a downstream config reader is meant to call, +with ``CALLER_ANNOTATION_POLICY``. The strict caller policy governs +the new ``pipeline_annotations`` input; the lenient document policy governs +every compiled or hand-authored pipeline and is deliberately unchanged +(scalar-or-null values, no key rules), so legacy YAML that validates today +still validates. The document check is also the backstop no path can bypass — +annotations that never went through the caller entry point still meet it. The +compiler contributes only orchestration and its precise error type. +""" + +from __future__ import annotations + +import shutil +import textwrap +from pathlib import Path +from unittest.mock import MagicMock + +import pytest +import yaml + +from tangle_cli import pipeline_compiler as pipeline_compiler_module +from tangle_cli.pipeline_compiler import PipelineCompiler, compile_pipeline +from tangle_cli.pipelines import compile_pipeline_file +from tangle_cli.python_pipeline.errors import ( + CompileError, + InvalidPipelineAnnotationsError, +) +from tangle_cli.schema_validation import ( + CALLER_ANNOTATION_POLICY, + DOCUMENT_ANNOTATION_POLICY, + RESERVED_ANNOTATION_KEY_PREFIX, + SchemaValidationError, + check_annotations, + validate_dehydrated_pipeline, +) + +FIXTURES = Path(__file__).parent / "fixtures" / "python_pipeline" + + +def _validate_caller(annotations): + """Validate exactly as the compiler does. + + The rules and the entry point live in the shared validation layer; the + compiler contributes only the caller-facing error type, so tests call the + shared API the same way a downstream config reader would. + """ + return check_annotations( + annotations, + policy=CALLER_ANNOTATION_POLICY, + error_cls=InvalidPipelineAnnotationsError, + ) + +#: A ``@task`` pipeline (hermetic — no external component YAML to colocate) +#: whose ``@pipeline`` decorator arguments are substituted per fixture. +_SOURCE_TEMPLATE = textwrap.dedent( + ''' + from tangle_cli.python_pipeline import In, Out, pipeline, task + + + @task(image="python:3.12") + def greet(greeting: str = "hi"): + """Write a greeting. + + Metadata: + Name: Greet + """ + print(greeting) + + + @pipeline(__DECORATOR_ARGS__) + def annotated_pipeline(seed: In[str]) -> Out[str]: + run_greet = greet(wait_for=seed) + return run_greet + ''' +) + +#: Source that already declares annotations, so merge/collision is observable. +_ANNOTATED_SOURCE = _SOURCE_TEMPLATE.replace( + "__DECORATOR_ARGS__", + '"Annotated Pipeline", ' + 'annotations={"author": "source-author", "version": "1.0"}', +) + +#: The same pipeline with NO source annotations, so the caller's mapping is the +#: only thing that can create the ``metadata`` block. +_UNANNOTATED_SOURCE = _SOURCE_TEMPLATE.replace( + "__DECORATOR_ARGS__", '"Annotated Pipeline"' +) + + +def _write(tmp_path: Path, source: str, name: str = "pipeline.py") -> Path: + src_dir = tmp_path / "src" + src_dir.mkdir(parents=True, exist_ok=True) + script = src_dir / name + script.write_text(source, encoding="utf-8") + return script + + +def _annotations_of(path: Path) -> dict: + data = yaml.safe_load(path.read_text(encoding="utf-8")) + return data.get("metadata", {}).get("annotations", {}) + + +def _bundle_bytes(root: Path) -> dict[str, bytes]: + """Every file the compile wrote under ``root``'s directory, by name.""" + return { + str(p.relative_to(root.parent)): p.read_bytes() + for p in sorted(root.parent.rglob("*")) + if p.is_file() + } + + +# --------------------------------------------------------------------------- +# Merge semantics + + +def test_caller_annotations_merge_per_key_and_win_on_collision(tmp_path): + script = _write(tmp_path, _ANNOTATED_SOURCE) + out = tmp_path / "out" / "compiled.yaml" + + compile_pipeline( + script, + out, + pipeline_annotations={"version": "2.0", "environment": "staging"}, + ) + + annotations = _annotations_of(out) + # Collision resolves to the caller; the untouched source key survives; the + # caller's new key is added. Source-declared order comes first. + assert annotations == { + "author": "source-author", + "version": "2.0", + "environment": "staging", + } + assert list(annotations) == ["author", "version", "environment"] + + +def test_caller_annotations_create_the_metadata_block_when_source_has_none(tmp_path): + script = _write(tmp_path, _UNANNOTATED_SOURCE) + out = tmp_path / "out" / "compiled.yaml" + + compile_pipeline(script, out, pipeline_annotations={"environment": "staging"}) + + assert _annotations_of(out) == {"environment": "staging"} + + +def test_the_caller_mapping_is_copied_not_aliased(tmp_path): + """Mutating the caller's mapping after the call cannot reach the output.""" + script = _write(tmp_path, _UNANNOTATED_SOURCE) + out = tmp_path / "out" / "compiled.yaml" + supplied = {"environment": "staging"} + + compile_pipeline(script, out, pipeline_annotations=supplied) + supplied["environment"] = "production" + supplied["late"] = "addition" + + assert _annotations_of(out) == {"environment": "staging"} + + +@pytest.mark.parametrize("supplied", [None, {}], ids=["absent", "empty"]) +def test_absent_or_empty_annotations_are_a_byte_identical_no_op(supplied, tmp_path): + """``{}`` is a no-op, NOT a clear of the source block. + + Both compiles write to the SAME paths, so byte identity is over the whole + bundle (root + ``@task`` sidecar) rather than over one file. + """ + script = _write(tmp_path, _ANNOTATED_SOURCE) + out = tmp_path / "out" / "compiled.yaml" + + compile_pipeline(script, out) + baseline = _bundle_bytes(out) + + compile_pipeline(script, out, pipeline_annotations=supplied) + + assert _bundle_bytes(out) == baseline + assert _annotations_of(out) == {"author": "source-author", "version": "1.0"} + + +# --------------------------------------------------------------------------- +# Root only: children inherit nothing, so their digests cannot move. + + +def test_a_subpipeline_child_neither_inherits_nor_changes_its_sidecar(tmp_path): + """The child sidecar's NAME and BYTES are identical with and without the + caller's annotations — the claim that component digests are untouched. + + The name is hashed over compile IDENTITY (``PipelineCompileKey``), so this + also pins that the value stays out of ``overrides_fingerprint``; the bytes + pin that no annotation leaked into the child graph. + """ + script = FIXTURES / "subpipeline_pipeline.py" + out = tmp_path / "out" / "compiled.yaml" + + baseline = compile_pipeline(script, out, pipeline_name="Parent Pipeline") + child_baseline = baseline.subgraph_paths[0] + child_baseline_bytes = child_baseline.read_bytes() + + annotated = compile_pipeline( + script, + out, + pipeline_name="Parent Pipeline", + pipeline_annotations={"environment": "staging"}, + ) + + assert len(annotated.subgraph_paths) == 1 + child = annotated.subgraph_paths[0] + assert child.name == child_baseline.name + assert child.read_bytes() == child_baseline_bytes + # Nothing landed in the child graph... + assert _annotations_of(child) == {} + # ...and everything landed in the root. + assert _annotations_of(out) == {"environment": "staging"} + + +# --------------------------------------------------------------------------- +# Hydration preserves root metadata. + + +def test_root_annotations_survive_hydration(tmp_path): + from tangle_cli.pipeline_hydrator import PipelineHydrator + + script = _write(tmp_path, _ANNOTATED_SOURCE) + # Compiled NEXT TO the source: hydration refuses to execute a + # ``local_from_python`` component whose Python file is not colocated with + # the bundle (or otherwise allowlisted). + out = script.parent / "compiled.yaml" + compile_pipeline( + script, out, pipeline_annotations={"version": "2.0", "environment": "staging"} + ) + + hydrated = PipelineHydrator(client=MagicMock()).hydrate_file(out) + + assert hydrated.data["metadata"]["annotations"] == { + "author": "source-author", + "version": "2.0", + "environment": "staging", + } + + +# --------------------------------------------------------------------------- +# Entry points: the handler and the facade carry the kwarg through. + + +def test_the_handler_and_the_facade_accept_the_kwarg(tmp_path): + script = _write(tmp_path, _UNANNOTATED_SOURCE) + + handler_out = tmp_path / "handler" / "compiled.yaml" + PipelineCompiler().compile_file( + script, handler_out, pipeline_annotations={"environment": "staging"} + ) + assert _annotations_of(handler_out) == {"environment": "staging"} + + facade_out = tmp_path / "facade" / "compiled.yaml" + compile_pipeline_file( + script, facade_out, pipeline_annotations={"environment": "production"} + ) + assert _annotations_of(facade_out) == {"environment": "production"} + + +# --------------------------------------------------------------------------- +# Validation: hostile mappings are refused, and no diagnostic echoes a value. + +#: Used as every rejected VALUE so one assertion can prove the message is +#: value-free regardless of which rule fired. +SECRET = "s3cret-annotation-payload" + + +@pytest.mark.parametrize( + ("supplied", "expected_fragment"), + [ + pytest.param( + [("environment", "staging")], + "must be a mapping of string keys to string values; got list", + id="non-mapping", + ), + pytest.param( + {1: SECRET}, + "keys must be strings; got a key of type int", + id="non-string-key", + ), + pytest.param( + {"": SECRET}, + "keys must not be empty", + id="empty-key", + ), + pytest.param( + {"environment": {"nested": SECRET}}, + "value for key 'environment' must be a string; got dict", + id="mapping-value", + ), + pytest.param( + {"enabled": True}, + "value for key 'enabled' must be a string; got bool", + id="bool-value", + ), + pytest.param( + {"replicas": 3}, + "value for key 'replicas' must be a string; got int", + id="int-value", + ), + pytest.param( + {"environment": None}, + "value for key 'environment' must be a string; got NoneType", + id="none-value", + ), + pytest.param( + {"system/owner": "platform"}, + "key 'system/owner' uses the reserved 'system/' prefix", + id="reserved-prefix", + ), + pytest.param( + {"environment": "{{ " + SECRET + " }}"}, + "value for key 'environment' contains the template delimiter '{{'", + id="jinja-expression-value", + ), + pytest.param( + {"environment": "{% if " + SECRET + " %}x{% endif %}"}, + "value for key 'environment' contains the template delimiter '{%'", + id="jinja-statement-value", + ), + pytest.param( + {"environment": "{# " + SECRET + " #}"}, + "value for key 'environment' contains the template delimiter '{#'", + id="jinja-comment-value", + ), + pytest.param( + {"{{ key }}": SECRET}, + "pipeline_annotations key '{{ key }}' contains the template " + "delimiter '{{'", + id="delimiter-in-key", + ), + ], +) +def test_a_hostile_annotations_mapping_is_refused_without_echoing_a_value( + supplied, expected_fragment +): + with pytest.raises(InvalidPipelineAnnotationsError) as exc: + _validate_caller(supplied) + + message = str(exc.value) + assert expected_fragment in message + # Untrusted input: the offending VALUE never reaches the diagnostic. + assert SECRET not in message + + +def test_the_annotations_error_is_a_compile_error(): + """Downstream callers may catch it precisely OR as a CompileError.""" + with pytest.raises(CompileError): + _validate_caller({"": "x"}) + assert issubclass(InvalidPipelineAnnotationsError, CompileError) + + +def test_validation_accepts_an_empty_value_and_a_lone_brace(): + """Only the three template delimiters are refused, not any brace.""" + assert _validate_caller({"environment": ""}) == {"environment": ""} + assert _validate_caller({"shape": "{json}"}) == {"shape": "{json}"} + + +@pytest.mark.parametrize("supplied", [None, {}], ids=["absent", "empty"]) +def test_validation_normalizes_nothing_to_an_empty_dict(supplied): + assert _validate_caller(supplied) == {} + + +def test_hostile_annotations_fail_before_the_script_is_read(tmp_path): + """Parse-time validation: a bad mapping is refused before the compiler + imports anything or writes anything, so the failure is about the + annotations rather than about whatever else the compile would hit.""" + missing = tmp_path / "src" / "does-not-exist.py" + out = tmp_path / "out" / "compiled.yaml" + + with pytest.raises(InvalidPipelineAnnotationsError) as exc: + compile_pipeline(missing, out, pipeline_annotations={"": SECRET}) + + assert "keys must not be empty" in str(exc.value) + assert not out.exists() + + +def test_a_delimiter_value_is_refused_by_the_compiler_not_the_output_scan(tmp_path): + """The compiled bundle must be fully hydrated; the caller-facing message + names the annotation key rather than a JSON path in the emitted YAML.""" + script = _write(tmp_path, _UNANNOTATED_SOURCE) + out = tmp_path / "out" / "compiled.yaml" + + with pytest.raises(InvalidPipelineAnnotationsError) as exc: + compile_pipeline( + script, out, pipeline_annotations={"environment": "{{ env_name }}"} + ) + + message = str(exc.value) + assert "pipeline_annotations value for key 'environment'" in message + assert "compiled pipeline output must contain no template delimiters" not in message + assert not out.exists() + + +# --------------------------------------------------------------------------- +# One shared policy: both entry points go through check_annotations, and the +# lenient document rules are unchanged for legacy YAML. + + +def test_the_compiler_delegates_to_the_shared_check(monkeypatch, tmp_path): + """The compiler owns no annotation rules: it calls the shared check with + the caller policy and only supplies its own precise error type.""" + seen = {} + + def _spy(annotations, *, policy=None, error_cls=None): + seen["annotations"] = annotations + seen["policy"] = policy + seen["error_cls"] = error_cls + return {} + + monkeypatch.setattr(pipeline_compiler_module, "check_annotations", _spy) + script = _write(tmp_path, _UNANNOTATED_SOURCE) + + compile_pipeline( + script, + tmp_path / "out" / "compiled.yaml", + pipeline_annotations={"environment": "staging"}, + ) + + assert seen["annotations"] == {"environment": "staging"} + assert seen["policy"] is CALLER_ANNOTATION_POLICY + assert seen["error_cls"] is InvalidPipelineAnnotationsError + + +def test_the_caller_policy_is_reusable_with_the_layer_default_error(): + """Same rules and messages whoever calls it. The DEFAULT error type is the + validation layer's own, so a caller opts into a precise one.""" + with pytest.raises(SchemaValidationError) as exc: + check_annotations({"replicas": 3}, policy=CALLER_ANNOTATION_POLICY) + + assert "pipeline_annotations value for key 'replicas' must be a string" in str( + exc.value + ) + assert check_annotations( + {"environment": "staging"}, policy=CALLER_ANNOTATION_POLICY + ) == {"environment": "staging"} + + +def test_the_reserved_prefix_has_one_definition(): + """The constant lives in the validation layer; the compiler keeps no copy.""" + assert RESERVED_ANNOTATION_KEY_PREFIX == "system/" + assert not hasattr(pipeline_compiler_module, "RESERVED_ANNOTATION_PREFIX") + + +@pytest.mark.parametrize( + "value", + ["text", 3, 1.5, True, None], + ids=["str", "int", "float", "bool", "null"], +) +def test_a_document_still_accepts_every_legacy_scalar_annotation(value): + """COMPATIBILITY: hand-authored / legacy YAML carrying a non-string scalar + annotation value must keep validating. The strict ``str -> str`` rule + applies ONLY to the caller-supplied input surface.""" + document = { + "name": "Legacy Pipeline", + "metadata": {"annotations": {"version": value, "": value}}, + "implementation": { + "graph": { + "tasks": { + "Only Task": { + "componentRef": {"url": "https://example.test/c.yaml"} + } + } + } + }, + } + + validate_dehydrated_pipeline(document) + # The same value is refused on the caller surface — the asymmetry is the + # point, and it is policy-driven rather than duplicated logic. + if not isinstance(value, str): + with pytest.raises(InvalidPipelineAnnotationsError): + _validate_caller({"version": value}) + + +def test_a_document_still_rejects_a_non_scalar_annotation_with_its_own_message(): + with pytest.raises(SchemaValidationError) as exc: + check_annotations( + {"version": {"nested": SECRET}}, policy=DOCUMENT_ANNOTATION_POLICY + ) + + message = str(exc.value) + assert ( + "metadata.annotations['version'] must be a scalar (str/number/bool) " + "or null, got 'dict'." in message + ) + assert SECRET not in message + + +@pytest.mark.parametrize( + "annotations", + [ + {"system/owner": "platform"}, + {"": "empty-key"}, + {"environment": "{{ env }}"}, + ], + ids=["reserved-prefix", "empty-key", "delimiter"], +) +def test_the_document_policy_is_not_tightened_by_the_caller_rules(annotations): + """Rules that exist only for caller input must NOT start rejecting + documents — the delimiter case still fails the OUTPUT scan, which is a + pre-existing guard, but the annotations check itself accepts it.""" + assert check_annotations(annotations, policy=DOCUMENT_ANNOTATION_POLICY) == dict( + annotations + ) + with pytest.raises(InvalidPipelineAnnotationsError): + _validate_caller(annotations) + + +def test_source_annotations_cannot_bypass_document_validation(tmp_path): + """An annotation that never passes through the caller entry point — here + one authored in ``@pipeline(annotations=...)`` — is still validated before + anything is written, so no path reaches a written bundle unchecked. + + Which layer refuses it (JSON-Schema or the shared semantic policy, whose + accepted value sets are identical by construction) is not the point and is + deliberately not asserted; that the annotation is named, and that nothing + is written, is. The compiler re-raises validation failures as + :class:`CompileError`, so this is NOT the caller-input error type. + """ + script = _write( + tmp_path, + _SOURCE_TEMPLATE.replace( + "__DECORATOR_ARGS__", + '"Annotated Pipeline", annotations={"owner": {"nested": "team"}}', + ), + ) + out = tmp_path / "out" / "compiled.yaml" + + with pytest.raises(CompileError) as exc: + compile_pipeline(script, out) + + assert not isinstance(exc.value, InvalidPipelineAnnotationsError) + assert "annotations" in str(exc.value) + assert "owner" in str(exc.value) + assert not out.exists() + + +def test_a_merged_caller_annotation_also_meets_the_document_check(tmp_path): + """The merged result is validated as a document like any other, so the + caller path is strict input validation layered ON TOP of the backstop, + not a replacement for it.""" + script = _write(tmp_path, _UNANNOTATED_SOURCE) + out = tmp_path / "out" / "compiled.yaml" + + compile_pipeline(script, out, pipeline_annotations={"environment": "staging"}) + + document = yaml.safe_load(out.read_text(encoding="utf-8")) + validate_dehydrated_pipeline(document) + + +# --------------------------------------------------------------------------- +# The existing fixtures still compile unchanged. + + +def test_an_ordinary_compile_is_unaffected(tmp_path): + out = tmp_path / "compiled.yaml" + out.parent.mkdir(parents=True, exist_ok=True) + shutil.copy(FIXTURES / "noop.yaml", out.parent / "noop.yaml") + + compile_pipeline(FIXTURES / "pipeline.py", out) + + data = yaml.safe_load(out.read_text(encoding="utf-8")) + assert "metadata" not in data diff --git a/uv.lock b/uv.lock index a6282e5..6a0724f 100644 --- a/uv.lock +++ b/uv.lock @@ -2083,7 +2083,7 @@ requires-dist = [{ name = "pydantic", specifier = ">=2.0" }] [[package]] name = "tangle-cli" -version = "0.1.18" +version = "0.1.19" source = { editable = "." } dependencies = [ { name = "cloud-pipelines" },