diff --git a/README.md b/README.md index 888e84d..9af821c 100644 --- a/README.md +++ b/README.md @@ -654,6 +654,31 @@ Semantics: 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. +##### Root pipeline labels + +`@pipeline(labels={...})` writes the compiled pipeline's root `metadata.labels` block — a block separate from `annotations`: + +```python +@pipeline( + "Search signals", + labels={"team": "discovery", "domain": "search-signals", "stage": "analysis"}, +) +def search_signals() -> Out[str]: + ... +``` + +Semantics: + +- **`str -> str`, and stricter than annotations.** Both the dehydrated and the pipeline schema type `metadata.labels` as `additionalProperties: {"type": "string"}`, whereas `metadata.annotations` also admits numbers, booleans and null. A non-mapping argument, a non-string key or value, an empty key, or a template delimiter (`{{`, `{%`, `{#`) raises `InvalidPipelineLabelsError` (a `CompileError`). Diagnostics name the key and the type and never echo a value. +- **The `system/` prefix is allowed here.** That prefix is reserved for Tangle's own *annotations*; nothing reserves a label prefix, so rejecting one would refuse a document both schemas accept. +- **Written before `annotations`** in the `metadata` block, whichever order the keywords were passed in. +- **Author key order is preserved** inside the block; it is not sorted. +- **Omitted, `{}` or `None` is a no-op**, byte for byte — a pipeline that does not use labels compiles exactly as before. +- **Root only.** A `subpipeline` child keeps exactly what its own `@pipeline` declared, so child sidecar names, bytes and component digests are unaffected. A child that wants labels declares its own. +- **Descriptive only.** Nothing in the CLI, the hydrator or the generated API client reads labels; they appear only in the schemas and in `MetadataSpec`. They are not part of compile identity or the overrides fingerprint, and cannot influence placement, routing, scheduling or run identity. + +There is no `pipeline_labels` compile keyword mirroring `pipeline_annotations`. Annotations got one because the varying part of that block (environment, owner) comes from a caller's per-environment config; no such need exists for labels today, and the rules are already shared, so adding one later is a small change. + 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. ##### Declaring graph inputs and outputs from the body diff --git a/packages/tangle-cli/src/tangle_cli/__init__.py b/packages/tangle-cli/src/tangle_cli/__init__.py index c6e0a11..d866dea 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.26" + __version__ = "0.1.27" __all__ = ["TangleDynamicDiscoveryClient", "__version__"] diff --git a/packages/tangle-cli/src/tangle_cli/python_pipeline/emit.py b/packages/tangle-cli/src/tangle_cli/python_pipeline/emit.py index 0d7f6ea..3df6950 100644 --- a/packages/tangle-cli/src/tangle_cli/python_pipeline/emit.py +++ b/packages/tangle-cli/src/tangle_cli/python_pipeline/emit.py @@ -82,9 +82,15 @@ def emit_pipeline(g: GraphBuilder) -> tuple[dict[str, Any], set[str]]: if g.description: out["description"] = g.description - if g.annotations: - # metadata.annotations preserves user-specified order. - out["metadata"] = {"annotations": dict(g.annotations)} + if g.labels or g.annotations: + # Both blocks preserve user-specified key order. ``labels`` is + # written first, matching the corpus majority. + metadata: dict[str, Any] = {} + if g.labels: + metadata["labels"] = dict(g.labels) + if g.annotations: + metadata["annotations"] = dict(g.annotations) + out["metadata"] = metadata if g.inputs: out["inputs"] = list(g.inputs) 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 86e12e9..a4eed4a 100644 --- a/packages/tangle-cli/src/tangle_cli/python_pipeline/errors.py +++ b/packages/tangle-cli/src/tangle_cli/python_pipeline/errors.py @@ -55,6 +55,15 @@ class InvalidGraphIoError(CompileError): """ +class InvalidPipelineLabelsError(CompileError): + """Raised on a malformed ``@pipeline(labels=...)`` mapping. + + Separate from :class:`InvalidPipelineAnnotationsError` so a caller can + tell which metadata block it got wrong. Messages name the key and the + type, never the value. + """ + + class InvalidPipelineAnnotationsError(CompileError): """Raised on a malformed caller-supplied ``pipeline_annotations`` mapping. diff --git a/packages/tangle-cli/src/tangle_cli/python_pipeline/graph.py b/packages/tangle-cli/src/tangle_cli/python_pipeline/graph.py index 1a2bc15..d682c41 100644 --- a/packages/tangle-cli/src/tangle_cli/python_pipeline/graph.py +++ b/packages/tangle-cli/src/tangle_cli/python_pipeline/graph.py @@ -67,6 +67,7 @@ class GraphBuilder: name: str description: str | None = None annotations: dict[str, str] = field(default_factory=dict) + labels: dict[str, str] = field(default_factory=dict) inputs: list[dict[str, Any]] = field(default_factory=list) outputs: list[dict[str, Any]] = field(default_factory=list) # MULTI-output map (Decision D): ``{output_name: EdgeRef}`` in field diff --git a/packages/tangle-cli/src/tangle_cli/python_pipeline/pipeline.py b/packages/tangle-cli/src/tangle_cli/python_pipeline/pipeline.py index 064860d..0f2803f 100644 --- a/packages/tangle-cli/src/tangle_cli/python_pipeline/pipeline.py +++ b/packages/tangle-cli/src/tangle_cli/python_pipeline/pipeline.py @@ -18,8 +18,10 @@ validate_flow_direction, ) +from tangle_cli.schema_validation import PIPELINE_LABELS_POLICY, check_annotations + from . import emit -from .errors import InvalidEditorLayoutError +from .errors import InvalidEditorLayoutError, InvalidPipelineLabelsError from .graph import GraphBuilder @@ -39,6 +41,9 @@ class PipelineFn: description: str | None = None config_path: str | None = None # path relative to caller_dir annotations: dict[str, Any] = field(default_factory=dict) + # ``metadata.labels`` — a separate block from annotations, and typed + # ``str -> str`` by both schemas. + labels: dict[str, str] = field(default_factory=dict) task_annotations: dict[str, Any] = field(default_factory=dict) caller_dir: Path | None = None # Convention for the single Out[T] slot's name. Defaults to the PoC @@ -135,6 +140,7 @@ def pipeline( description: str | None = None, config: str | None = None, annotations: dict[str, Any] | None = None, + labels: dict[str, str] | None = None, task_annotations: dict[str, Any] | None = None, flow_direction: str | None = None, output_name: str = "wait_for_output", @@ -156,6 +162,10 @@ def pipeline( time so ``--override key=value`` pairs can merge in. annotations: ``metadata.annotations`` block (e.g. ``version``, ``author``). + labels: ``metadata.labels`` block (e.g. ``team``, ``domain``, + ``stage``). A separate block from ``annotations``, restricted + by both schemas to string values. Descriptive only, and root + only — a ``subpipeline`` child declares its own. flow_direction: Editor rendering direction, written as the ``editor.flow-direction`` root annotation (``"left-to-right"`` / ``"top-to-bottom"``). Sugar over @@ -180,6 +190,11 @@ def decorator(fn: Callable[..., Any]) -> PipelineFn: caller_dir = None merged_annotations = dict(annotations or {}) + checked_labels = check_annotations( + labels, + policy=PIPELINE_LABELS_POLICY, + error_cls=InvalidPipelineLabelsError, + ) if flow_direction is not None: # Assigned after the mapping: the typed keyword wins, and an # existing key keeps its position in key order. @@ -193,6 +208,7 @@ def decorator(fn: Callable[..., Any]) -> PipelineFn: description=description, config_path=config, annotations=merged_annotations, + labels=checked_labels, task_annotations=dict(task_annotations or {}), caller_dir=caller_dir, output_name=output_name, diff --git a/packages/tangle-cli/src/tangle_cli/python_pipeline/trace.py b/packages/tangle-cli/src/tangle_cli/python_pipeline/trace.py index 01529df..c97027c 100644 --- a/packages/tangle-cli/src/tangle_cli/python_pipeline/trace.py +++ b/packages/tangle-cli/src/tangle_cli/python_pipeline/trace.py @@ -242,6 +242,7 @@ def trace_pipeline( name=pipeline_fn.name, description=pipeline_fn.description, annotations=dict(pipeline_fn.annotations), + labels=dict(pipeline_fn.labels), ) # AST pre-pass: needed by CallableRef.__call__ to derive task IDs # from LHS variable names. Stashed on the builder so the contextvar diff --git a/packages/tangle-cli/src/tangle_cli/schema_validation.py b/packages/tangle-cli/src/tangle_cli/schema_validation.py index ded305e..0011e7d 100644 --- a/packages/tangle-cli/src/tangle_cli/schema_validation.py +++ b/packages/tangle-cli/src/tangle_cli/schema_validation.py @@ -257,6 +257,22 @@ class AnnotationPolicy: reject_non_mapping=True, ) +#: Applied to a ``@pipeline(labels=...)`` mapping. Strict ``str -> str`` +#: because both schemas type ``metadata.labels`` as +#: ``additionalProperties: {"type": "string"}`` — narrower than +#: ``metadata.annotations``, which also admits numbers, booleans and null. +#: ``reject_reserved_key_prefix`` is deliberately OFF: ``system/`` is +#: documented as reserved for Tangle's own ANNOTATIONS, and nothing reserves +#: a label prefix, so enabling it would refuse a document the schema accepts. +PIPELINE_LABELS_POLICY = AnnotationPolicy( + label="labels", + require_string_values=True, + require_string_keys=True, + require_non_empty_keys=True, + reject_template_delimiters=True, + reject_non_mapping=True, +) + def check_annotations( annotations: Any, diff --git a/pyproject.toml b/pyproject.toml index ac4d1c1..1a8c40e 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "tangle-cli" -version = "0.1.26" +version = "0.1.27" 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 47f6124..ad022be 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.26" in metadata + assert "Version: 0.1.27" 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_labels.py b/tests/test_pipeline_labels.py new file mode 100644 index 0000000..fc650d5 --- /dev/null +++ b/tests/test_pipeline_labels.py @@ -0,0 +1,388 @@ +"""``@pipeline(labels=...)`` writes root ``metadata.labels``. + +``GraphBuilder`` carried only annotations, so the five corpus pipelines with +a root ``metadata.labels`` block could not be authored in Python. Both +schemas type ``labels`` as ``additionalProperties: {"type": "string"}`` — +narrower than ``annotations``, which also admits numbers, booleans and null +— so the value policy here is strict ``str -> str``. + +Labels are descriptive: nothing in the CLI, hydrator or client reads them. +""" + +from __future__ import annotations + +import re +import textwrap +from pathlib import Path + +import pytest +import yaml + +from tangle_cli.pipeline_compiler import compile_pipeline +from tangle_cli.python_pipeline import pipeline +from tangle_cli.python_pipeline.errors import InvalidPipelineLabelsError + +# The exact label blocks carried by the five corpus pipelines. +_CORPUS_LABELS = { + "search_signals": { + "team": "discovery", + "domain": "search-signals", + "stage": "analysis", + }, + "join_features": { + "team": "discovery", + "domain": "storefront-reranker", + "stage": "experiment", + }, + "l1_tangentable": { + "team": "discovery", + "stage": "experiment", + "domain": "shop-app-search-ranking", + "tangentable": "v3", + }, + "smoke_test": { + "team": "discovery", + "domain": "storefront-reranker", + "stage": "smoke-test", + }, +} + +_TEMPLATE = ''' +from tangle_cli.python_pipeline import Out, pipeline, task + + +@task(image="python:3.12") +def echo(value: str) -> str: + return value + + +@pipeline("Labelled"__DECORATOR_ARGS__) +def labelled() -> Out[str]: + run_it = echo.named("Echo")(value="x") + return run_it +''' + + +def _compile(tmp_path: Path, decorator_args: str, case: str) -> dict: + case_dir = tmp_path / case + case_dir.mkdir(parents=True, exist_ok=True) + script = case_dir / "pipeline.py" + script.write_text( + textwrap.dedent(_TEMPLATE).replace("__DECORATOR_ARGS__", decorator_args), + encoding="utf-8", + ) + out = case_dir / "compiled.yaml" + compile_pipeline(script, out) + return yaml.safe_load(out.read_text(encoding="utf-8")) + + +def _normalize(path: Path) -> str: + """Compiled text with the child sidecar's content hash masked, so two + compiles in different directories can be compared.""" + return re.sub(r"child-[0-9a-f]{8}", "child-HASH", path.read_text(encoding="utf-8")) + + +def _child_sidecars(case_dir: Path) -> list[Path]: + """Child subgraph YAMLs, excluding the sibling ``.components.yaml``.""" + return sorted( + path + for path in (case_dir / "compiled.subgraphs").glob("child-*.yaml") + if not path.name.endswith(".components.yaml") + ) + + +# ============================================================================ +# Emitted shape +# ============================================================================ + + +@pytest.mark.parametrize("case, labels", sorted(_CORPUS_LABELS.items())) +def test_each_corpus_label_block_round_trips(tmp_path, case, labels): + doc = _compile(tmp_path, f", labels={labels!r}", case) + + assert doc["metadata"] == {"labels": labels} + # Key order inside the block is the author's, not sorted. + assert list(doc["metadata"]["labels"]) == list(labels) + + +def test_labels_are_written_before_annotations(tmp_path): + """Corpus metadata blocks are 4:1 labels-first, so that is the canonical + order regardless of which keyword the author passed first.""" + doc = _compile( + tmp_path, + ', annotations={"version": "1.0"}, labels={"team": "discovery"}', + "both", + ) + + assert list(doc["metadata"]) == ["labels", "annotations"] + assert doc["metadata"] == { + "labels": {"team": "discovery"}, + "annotations": {"version": "1.0"}, + } + + +def test_labels_reach_the_yaml_text_as_plain_strings(tmp_path): + """Guards the rendered text, not just the parsed mapping.""" + case_dir = tmp_path / "text" + case_dir.mkdir() + script = case_dir / "pipeline.py" + script.write_text( + textwrap.dedent(_TEMPLATE).replace( + "__DECORATOR_ARGS__", ', labels={"team": "discovery", "stage": "analysis"}' + ), + encoding="utf-8", + ) + out = case_dir / "compiled.yaml" + compile_pipeline(script, out) + + text = out.read_text(encoding="utf-8") + assert "metadata:\n labels:\n team: discovery\n stage: analysis\n" in text + + +# ============================================================================ +# No labels changes nothing +# ============================================================================ + + +def test_a_pipeline_without_labels_is_byte_identical(tmp_path): + """The new block must not perturb any pipeline that does not use it.""" + without = _compile(tmp_path, "", "plain_a") + explicit_empty = _compile(tmp_path, ", labels={}", "plain_b") + explicit_none = _compile(tmp_path, ", labels=None", "plain_c") + + assert "metadata" not in without + assert without == explicit_empty == explicit_none + + a = (tmp_path / "plain_a" / "compiled.yaml").read_bytes() + b = (tmp_path / "plain_b" / "compiled.yaml").read_bytes() + c = (tmp_path / "plain_c" / "compiled.yaml").read_bytes() + assert a == b == c + + +def test_annotations_alone_still_emit_the_same_metadata_block(tmp_path): + doc = _compile(tmp_path, ', annotations={"version": "1.0"}', "ann_only") + + assert doc["metadata"] == {"annotations": {"version": "1.0"}} + + +# ============================================================================ +# Subpipeline +# ============================================================================ + + +def test_labels_do_not_inherit_into_a_subpipeline_child(tmp_path): + """Consistent with root annotations: a child keeps exactly what its own + ``@pipeline`` declared, so child sidecar bytes are unaffected.""" + case_dir = tmp_path / "subpipe" + case_dir.mkdir() + script = case_dir / "pipeline.py" + script.write_text( + textwrap.dedent( + ''' + from tangle_cli.python_pipeline import In, Out, pipeline, subpipeline, task + + + @task(image="python:3.12") + def echo(value: str) -> str: + return value + + + @pipeline("Child") + def child(seed: In[str]) -> Out[str]: + inner = echo.named("Inner")(value=seed) + return inner + + + @pipeline("Parent", labels={"team": "discovery"}) + def parent() -> Out[str]: + kid = subpipeline(child).named("Run child")(seed="x") + return kid.wait_for_output + ''' + ), + encoding="utf-8", + ) + out = case_dir / "compiled.yaml" + compile_pipeline(script, out, pipeline_name="parent") + + parent_doc = yaml.safe_load(out.read_text(encoding="utf-8")) + assert parent_doc["metadata"] == {"labels": {"team": "discovery"}} + + children = _child_sidecars(case_dir) + assert len(children) == 1 + child_doc = yaml.safe_load(children[0].read_text(encoding="utf-8")) + assert "metadata" not in child_doc + + # Stronger than "no metadata key": compile the same pair with no labels + # at all and compare. The child sidecar's name is a content hash that + # also folds in the output directory, so the two compiles necessarily + # live in different directories and the hash token is normalized away. + # Everything else must match, in BOTH documents. + unlabelled_parent, unlabelled_child = _compile_parent_without_labels(tmp_path) + assert _normalize(children[0]) == _normalize(unlabelled_child) + assert _normalize(out) == _normalize(unlabelled_parent).replace( + "name: Parent\n", "name: Parent\nmetadata:\n labels:\n team: discovery\n" + ) + + +def test_a_labelled_child_keeps_its_own_labels(tmp_path): + """The other direction: declaring labels on the child works, and does not + leak up to the parent.""" + case_dir = tmp_path / "subpipe_child" + case_dir.mkdir() + script = case_dir / "pipeline.py" + script.write_text( + textwrap.dedent( + ''' + from tangle_cli.python_pipeline import In, Out, pipeline, subpipeline, task + + + @task(image="python:3.12") + def echo(value: str) -> str: + return value + + + @pipeline("Child", labels={"stage": "analysis"}) + def child(seed: In[str]) -> Out[str]: + inner = echo.named("Inner")(value=seed) + return inner + + + @pipeline("Parent") + def parent() -> Out[str]: + kid = subpipeline(child).named("Run child")(seed="x") + return kid.wait_for_output + ''' + ), + encoding="utf-8", + ) + out = case_dir / "compiled.yaml" + compile_pipeline(script, out, pipeline_name="parent") + + parent_doc = yaml.safe_load(out.read_text(encoding="utf-8")) + assert "metadata" not in parent_doc + + child_doc = yaml.safe_load( + _child_sidecars(case_dir)[0].read_text(encoding="utf-8") + ) + assert child_doc["metadata"] == {"labels": {"stage": "analysis"}} + + +# ============================================================================ +# Validation +# ============================================================================ + + +def _reject(labels) -> str: + with pytest.raises(InvalidPipelineLabelsError) as excinfo: + + @pipeline("Rejected", labels=labels) + def _p(): # pragma: no cover - never traced + pass + + return str(excinfo.value) + + +def test_a_non_string_value_is_refused(): + """Stricter than annotations, which accept numbers, booleans and null: + both schemas type label values as string.""" + message = _reject({"replicas": 3}) + + assert "labels value for key 'replicas' must be a string" in message + assert "int" in message + + +def test_a_bool_value_is_refused(): + assert "must be a string" in _reject({"tangentable": True}) + + +def test_a_none_value_is_refused(): + assert "must be a string" in _reject({"stage": None}) + + +def test_a_non_string_key_is_refused(): + assert "keys must be strings" in _reject({3: "discovery"}) + + +def test_an_empty_key_is_refused(): + assert "must not be empty" in _reject({"": "discovery"}) + + +def test_a_non_mapping_is_refused(): + assert "must be a mapping" in _reject(["team=discovery"]) + + +def test_a_template_delimiter_is_refused(): + """Moves the failure off the whole-output delimiter scan and onto a + message that names the offending label key.""" + assert "team" in _reject({"team": "{{ team_name }}"}) + + +def test_the_system_prefix_is_allowed_on_labels(): + """``system/`` is reserved for Tangle's own ANNOTATIONS. Nothing reserves + a label prefix, so refusing it would reject a document both schemas + accept.""" + + @pipeline("Reserved-ish", labels={"system/owner": "discovery"}) + def _p(): # pragma: no cover - never traced + pass + + assert _p.labels == {"system/owner": "discovery"} + + +def test_diagnostics_never_echo_a_label_value(): + for labels in ( + {"team": 99999}, + {"team": {"nested": "super-secret-label"}}, + {"team": "{{ super-secret-label }}"}, + ): + message = _reject(labels) + assert "super-secret-label" not in message + assert "99999" not in message + + +def test_the_checked_mapping_is_copied(): + """A later mutation of the caller's dict must not reach the compile.""" + supplied = {"team": "discovery"} + + @pipeline("Copied", labels=supplied) + def _p(): # pragma: no cover - never traced + pass + + supplied["team"] = "changed" + assert _p.labels == {"team": "discovery"} + + +def _compile_parent_without_labels(tmp_path: Path) -> tuple[Path, Path]: + """Compile the same parent/child pair with no labels anywhere and return + ``(parent, child)`` paths, for the leak comparison above.""" + case_dir = tmp_path / "subpipe_baseline" + case_dir.mkdir() + script = case_dir / "pipeline.py" + script.write_text( + textwrap.dedent( + ''' + from tangle_cli.python_pipeline import In, Out, pipeline, subpipeline, task + + + @task(image="python:3.12") + def echo(value: str) -> str: + return value + + + @pipeline("Child") + def child(seed: In[str]) -> Out[str]: + inner = echo.named("Inner")(value=seed) + return inner + + + @pipeline("Parent") + def parent() -> Out[str]: + kid = subpipeline(child).named("Run child")(seed="x") + return kid.wait_for_output + ''' + ), + encoding="utf-8", + ) + out = case_dir / "compiled.yaml" + compile_pipeline(script, out, pipeline_name="parent") + return out, _child_sidecars(case_dir)[0] diff --git a/uv.lock b/uv.lock index b0f9507..6b51d98 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.26" +version = "0.1.27" source = { editable = "." } dependencies = [ { name = "cloud-pipelines" },