Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 28 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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=`:
Expand Down
2 changes: 1 addition & 1 deletion packages/tangle-cli/src/tangle_cli/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,6 @@
try:
__version__ = metadata_version("tangle-cli")
except PackageNotFoundError:
__version__ = "0.1.18"
__version__ = "0.1.19"

__all__ = ["TangleDynamicDiscoveryClient", "__version__"]
40 changes: 38 additions & 2 deletions packages/tangle-cli/src/tangle_cli/pipeline_compiler.py
Original file line number Diff line number Diff line change
Expand Up @@ -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


Expand Down Expand Up @@ -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``.

Expand Down Expand Up @@ -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
Expand All @@ -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()
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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``.

Expand All @@ -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:
Expand Down
11 changes: 11 additions & 0 deletions packages/tangle-cli/src/tangle_cli/pipelines.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
11 changes: 11 additions & 0 deletions packages/tangle-cli/src/tangle_cli/python_pipeline/errors.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
"""
Loading
Loading