From 85a980a090547a761621681ce10de936f1177f82 Mon Sep 17 00:00:00 2001 From: Volv G Date: Fri, 2 Oct 2026 18:02:06 -0700 Subject: [PATCH] Add compile-time @Layout() graph auto-layout seam (v0.1.29) `@Layout()` / `@Layout("name", recursive=...)` marks a @pipeline graph definition for auto-layout. When the caller passes a `layout_transform` to `compile_pipeline` (also forwarded by `PipelineCompiler.compile_file` and `pipelines.compile_pipeline_file`), every covered graph is laid out post-order before validation and writing, so the root and subgraph YAML carry `editor.position` annotations. Tangle CLI ships no algorithm. - Coverage: an explicit declaration wins; `recursive=True` (default) covers undecorated descendants, `recursive=False` only its own graph; `@Layout()` resets to the transform's default. - Shared children compile per (compile key, layout policy): different inherited policies yield distinct sidecars, the same policy dedups. Uncovered children keep legacy filenames; cycle detection ignores layout. - The transform receives a private copy plus a `GraphLayoutContext` (governing Layout, first-occurrence path, task interfaces, artifact dir). A guard allows only `editor.position` changes on graph tasks and top-level inputs/outputs; failures write nothing. - `TaskInterface` is exact for @task/subpipeline tasks and approximate (supplied arguments + consumed outputs, incl. isEnabled) for opaque refs; nothing is fetched remotely. - Without a transform, or without @Layout(), output is byte-identical (a warning notes each unapplied @Layout()). --- README.md | 31 + .../tangle-cli/src/tangle_cli/__init__.py | 2 +- .../src/tangle_cli/pipeline_compiler.py | 368 ++++++- .../tangle-cli/src/tangle_cli/pipelines.py | 6 + .../tangle_cli/python_pipeline/__init__.py | 11 + .../python_pipeline/compiler_context.py | 17 +- .../src/tangle_cli/python_pipeline/errors.py | 10 + .../src/tangle_cli/python_pipeline/layout.py | 237 +++++ .../tangle_cli/python_pipeline/pipeline.py | 10 + pyproject.toml | 2 +- tests/test_layout_annotation.py | 904 ++++++++++++++++++ tests/test_packaging.py | 2 +- tests/test_python_pipeline_dsl.py | 5 + uv.lock | 2 +- 14 files changed, 1596 insertions(+), 11 deletions(-) create mode 100644 packages/tangle-cli/src/tangle_cli/python_pipeline/layout.py create mode 100644 tests/test_layout_annotation.py diff --git a/README.md b/README.md index 50db428..ad99c53 100644 --- a/README.md +++ b/README.md @@ -760,6 +760,37 @@ Graph **inputs** and **outputs** carry `editor.position` in the same way. An `In The key names, the serialized format, and the validation live in `tangle_cli.editor_layout` (`POSITION_ANNOTATION`, `FLOW_DIRECTION_ANNOTATION`, `FLOW_DIRECTIONS`, `position_annotation_value`, `validate_flow_direction`), which the CLI's own `pipelines layout` and the runner's auto-layout gate share, so a tool emitting layout cannot drift from what the editor reads. That module is stdlib-only and lives outside `python_pipeline` on purpose, so layout and submit paths do not pull in the authoring/codegen stack; it raises whatever `error_cls` the caller injects, and the authoring surfaces inject `InvalidEditorLayoutError`. +##### Compile-time auto-layout (`@Layout()`) + +`@Layout()` asks the compiler to auto-layout a graph, so that the compiled root and subgraph YAML already contain `editor.position` annotations. Tangle CLI ships no layout algorithm. The caller passes a `layout_transform` to `compile_pipeline(...)`, `PipelineCompiler.compile_file(...)` or `tangle_cli.pipelines.compile_pipeline_file(...)`; a downstream distribution typically installs one for all of its compile paths: + +```python +from tangle_cli.python_pipeline import Layout, pipeline, subpipeline + +@Layout("banded", recursive=False) # either order around @pipeline works +@pipeline("Judge") +def judge(...): ... + +@pipeline("Daily Pulse") +@Layout() # algorithm=None: the transform's default +def daily_pulse(...): + first = subpipeline(judge).named("First")(...) + ... +``` + +- **Coverage.** `@Layout(...)` covers its own graph. With `recursive=True` (the default) it also covers every undecorated descendant subgraph; with `recursive=False` it covers only its own graph. A descendant with its own `@Layout(...)` always uses that declaration, and its `recursive` flag governs further down. `@Layout()` therefore resets to the transform's default rather than inheriting an ancestor's algorithm. An undecorated root with decorated descendants lays out only those descendants. Only authored `@pipeline` graphs are laid out. External `ref(...)` / `@registered` components are never read remotely or modified. +- **The transform.** It is called as `transform(graph, context)` once per covered compiled artifact, post-order (subgraphs first). The call comes after the graph's references are final and before validation and writing. `graph` is a private deep copy of the compiled body. The function returns the body to write, which may differ only in `editor.position` string values on graph tasks and on top-level `inputs`/`outputs`. Both creating and replacing positions are allowed, and an explicit `.with_position(...)` is overwritten. Any other difference, a non-string position or a non-dict return value raises `CompileError`, and nothing is written. A `CompileError` raised by the transform propagates as-is; any other exception is wrapped. The `GraphLayoutContext` carries: + - `layout`: the governing `Layout`; + - `path`: the task-ID path of the first occurrence, with the root `()`; + - `pipeline_name`; + - `task_interfaces`: every task ID → `TaskInterface(inputs, outputs, approximate)` with ordered port names. `@task` and `subpipeline` tasks get their exact interface, extended with any other observed names. Opaque `ref` / `@registered` tasks get observed names only, with `approximate=True`: every argument the task supplies (edges and filled literal values alike) as inputs, and the outputs consumed elsewhere in the graph as outputs, so their cards may be smaller than the real component. Nothing is hydrated remotely, and interface data is never written back; + - `artifact_dir`. + + The transform must be deterministic, because occurrences that share an artifact share its result. +- **Variants and filenames.** A shared child is still compiled once per definition, config and layout policy. The same child reached under two different policies (for example `A` from one parent and `B` from another) becomes two correctly laid-out sidecars with distinct filenames; the same policy still dedups, diamonds included. Uncovered children keep their existing sidecar filenames. Cycle detection ignores layout. +- **No transform, no change.** Without a `layout_transform`, `@Layout()` changes nothing in the output: bytes, filenames and keys are identical to an undecorated compile, and `CompileResult.warnings` notes each decorated pipeline. Without `@Layout()`, an installed transform is never called. +- **Validated structurally, without echoing values.** `algorithm` must be `None` or a non-empty `str` (positional or keyword); which names are supported is the transform's decision. `recursive` must be a `bool`. A bare `@Layout` (without parentheses), a second `@Layout()` on the same definition, or a non-pipeline target (a `@task`, a `subpipeline(...)` handle, a class, a partial, a `PipelineFn` subclass, …) raises `InvalidLayoutError` (a `CompileError`). `Layout` objects are immutable, and the request stays attached to the `PipelineFn`, so an imported or cached child keeps it. + ##### 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 b8a2fd4..eaedcd3 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.28" + __version__ = "0.1.29" __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 1c4826d..34257b6 100644 --- a/packages/tangle-cli/src/tangle_cli/pipeline_compiler.py +++ b/packages/tangle-cli/src/tangle_cli/pipeline_compiler.py @@ -36,6 +36,7 @@ from __future__ import annotations +import copy import hashlib import importlib.util import inspect @@ -43,7 +44,9 @@ import os import re import sys +import types import uuid +import warnings as _warnings from collections.abc import Iterator, Mapping from contextlib import contextmanager from dataclasses import dataclass, field @@ -54,6 +57,7 @@ from .authenticated_identity import ME from .component_from_func import build_unwrapped_inputs_schema +from .editor_layout import POSITION_ANNOTATION from .handler import TangleCliHandler from .python_pipeline.cfg import ( Cfg, @@ -73,6 +77,12 @@ ) from .python_pipeline.emit import _TASK_URL_PLACEHOLDER, emit_pipeline from .python_pipeline.errors import CompileError, InvalidPipelineAnnotationsError +from .python_pipeline.layout import ( + GraphLayoutContext, + GraphLayoutTransform, + Layout, + TaskInterface, +) from .python_pipeline.pipeline import PipelineFn from .python_pipeline.ref import CallableRef from .python_pipeline.registered import _REGISTERED_URL_PLACEHOLDER @@ -154,6 +164,7 @@ def compile_pipeline( emit_components_sidecar: bool = True, image_overrides: Mapping[str, str] | None = None, pipeline_annotations: Mapping[str, str] | None = None, + layout_transform: GraphLayoutTransform | None = None, ) -> CompileResult: """Compile ``script`` to a single pipeline YAML at ``output``. @@ -199,6 +210,13 @@ def compile_pipeline( :data:`tangle_cli.schema_validation.CALLER_ANNOTATION_POLICY` for the accepted shape; a malformed mapping raises :class:`~tangle_cli.python_pipeline.errors.InvalidPipelineAnnotationsError`. + layout_transform: Optional :class:`GraphLayoutTransform` that lays out + every graph covered by ``@Layout()`` before validation and writing, + so the written root/sidecar YAML carries its ``editor.position`` + annotations. See :func:`_layout_policy_for` for coverage and + :func:`_apply_layout_transform` for the call contract. ``None`` + (the default) leaves output byte-identical to an undecorated + compile and adds one warning per decorated pipeline. Returns: A :class:`CompileResult`. ``components_path`` is the sidecar path @@ -276,6 +294,7 @@ def compile_pipeline( source_dirs=purge_dirs, image_overrides=image_overrides, pipeline_annotations=root_annotations, + layout_transform=layout_transform, ) # 5. Compile the root (and, recursively, all children) into in-memory @@ -289,6 +308,8 @@ def compile_pipeline( overrides, is_root=True, base_dir=script_path.parent, + layout_policy=_layout_policy_for(pipeline_fn, None, ctx), + occurrence_path=(), ) # 6. The full bundle: the root plus every deduped child in the registry. @@ -320,6 +341,308 @@ def compile_pipeline( _purge_bundle_local_modules(purge_dirs) +# --------------------------------------------------------------------------- +# Compile-time graph layout (``@Layout()``) + + +@dataclass(frozen=True) +class _LayoutPolicy: + """Layout coverage of one graph occurrence. + + ``apply`` is the :class:`Layout` this graph is laid out with (``None``: + not laid out); ``carry`` is the one its undecorated subgraphs inherit + (``None``: they inherit nothing). Both always ``None`` without a + transform, which keeps plain compiles identical. + """ + + apply: Layout | None = None + carry: Layout | None = None + + def variant(self) -> str | None: + """Canonical artifact-variant fingerprint, ``None`` when uncovered. + + Folded into a child's dedup key and sidecar filename, so a child laid + out under two different policies becomes two artifacts while the same + policy still dedups. ``None`` keeps the pre-layout key and filename. + """ + if self.apply is None and self.carry is None: + return None + + def encode(layout: Layout | None) -> dict[str, Any] | None: + if layout is None: + return None + return {"algorithm": layout.algorithm, "recursive": layout.recursive} + + return json.dumps( + {"apply": encode(self.apply), "carry": encode(self.carry)}, + sort_keys=True, + separators=(",", ":"), + ) + + +_NO_LAYOUT = _LayoutPolicy() + + +def _layout_policy_for( + pipeline_fn: PipelineFn, + inherited: Layout | None, + ctx: CompileContext, +) -> _LayoutPolicy: + """Resolve the layout coverage of one occurrence of ``pipeline_fn``. + + The graph's own ``@Layout()`` wins over anything inherited, and its + ``recursive`` flag decides whether undecorated descendants inherit it. + An undecorated graph is laid out with, and passes down, whatever its + nearest recursive ancestor declared. Without a transform nothing is + covered, and each decorated pipeline produces one warning. + """ + explicit = pipeline_fn.layout + if ctx.layout_transform is None: + if explicit is not None and pipeline_fn.name not in ctx.layout_warned: + ctx.layout_warned.add(pipeline_fn.name) + ctx.warnings.append( + f"pipeline {pipeline_fn.name!r} has @Layout() but no layout transform " + "is installed for this compile, so no auto-layout positions were " + "written. Compile through a tool that provides a layout transform." + ) + return _NO_LAYOUT + if explicit is not None: + return _LayoutPolicy(apply=explicit, carry=explicit if explicit.recursive else None) + if inherited is not None: + return _LayoutPolicy(apply=inherited, carry=inherited) + return _NO_LAYOUT + + +def _variant_hash8(key: PipelineCompileKey, variant: str | None) -> str: + """Sidecar filename hash: the plain key hash unless a layout variant applies.""" + if variant is None: + return key.hash8() + blob = json.dumps({**key.canonical_fields(), "layout": variant}, sort_keys=True, separators=(",", ":")) + return hashlib.sha256(blob.encode("utf-8")).hexdigest()[:8] + + +def _task_interface_for_ref(ref: CallableRef, unwrapped: Mapping[str, Any] | None) -> TaskInterface | None: + """Ordered component interface of an ``@task`` call, or ``None`` if unknown. + + Built from the traced function with the same ``extract_interface`` the + hydrator later generates the component from (best effort: geometry only). + """ + fn = vars(ref).get("__wrapped__") + if not callable(fn): + return None + try: + from .component_from_func import extract_interface + + with _warnings.catch_warnings(): + _warnings.simplefilter("ignore") + spec = extract_interface(fn, {}, unwrapped_inputs=dict(unwrapped) if unwrapped else None) + except Exception: # noqa: BLE001 — geometry is best effort; never fail the compile here + return None + return TaskInterface( + inputs=tuple(p.yaml_name for p in spec.inputs), + outputs=tuple(p.yaml_name for p in spec.all_outputs), + ) + + +def _graph_io_names(body: Mapping[str, Any], key: str) -> tuple[str, ...]: + specs = body.get(key) or [] + return tuple(s["name"] for s in specs if isinstance(s, dict) and isinstance(s.get("name"), str)) + + +def _observed_task_ports(body: Mapping[str, Any]) -> dict[str, tuple[list[str], list[str]]]: + """Task ID -> ``(input names, output names)`` observed in ``body``. + + Inputs are EVERY argument name the task supplies, in argument order: + edges (``taskOutput`` / ``graphInput``) and filled literal values alike. + Outputs are the distinct ``taskOutput`` names other tasks consume, as + arguments or in their ``isEnabled`` condition, or that ``outputValues`` + expose, in first-seen order. ``isEnabled`` is not an input port. + """ + graph = body.get("implementation", {}).get("graph", {}) or {} + tasks = graph.get("tasks", {}) or {} + observed: dict[str, tuple[list[str], list[str]]] = {} + for task_id, task_spec in tasks.items(): + arguments = task_spec.get("arguments") if isinstance(task_spec, dict) else None + names = [name for name in arguments if isinstance(name, str)] if isinstance(arguments, dict) else [] + observed[task_id] = (names, []) + + def consume(value: Any) -> None: + task_output = value.get("taskOutput") if isinstance(value, dict) else None + if not isinstance(task_output, dict): + return + target = observed.get(task_output.get("taskId")) # type: ignore[arg-type] + name = task_output.get("outputName") + if target is not None and isinstance(name, str) and name not in target[1]: + target[1].append(name) + + def consume_nested(value: Any) -> None: + # ``isEnabled`` may be a bare reference or a nested predicate. + if isinstance(value, dict): + consume(value) + for nested in value.values(): + consume_nested(nested) + elif isinstance(value, list): + for nested in value: + consume_nested(nested) + + for task_spec in tasks.values(): + if not isinstance(task_spec, dict): + continue + arguments = task_spec.get("arguments") + if isinstance(arguments, dict): + for value in arguments.values(): + consume(value) + consume_nested(task_spec.get("isEnabled")) + output_values = graph.get("outputValues") + if isinstance(output_values, dict): + for value in output_values.values(): + consume(value) + return observed + + +def _union_interface( + known: TaskInterface | None, + observed_inputs: list[str], + observed_outputs: list[str], +) -> TaskInterface: + """Known interface extended by any observed port it lacks; or the observed + ports alone, marked ``approximate``, when nothing is known.""" + if known is None: + return TaskInterface(inputs=tuple(observed_inputs), outputs=tuple(observed_outputs), approximate=True) + return TaskInterface( + inputs=known.inputs + tuple(n for n in observed_inputs if n not in known.inputs), + outputs=known.outputs + tuple(n for n in observed_outputs if n not in known.outputs), + ) + + +def _layout_task_interfaces( + artifact: SubgraphArtifact, + builder: Any, + sidecar_plan: TaskSidecarPlan | None, +) -> dict[str, TaskInterface]: + """Graph task ID -> :class:`TaskInterface`, in task order. + + Exact for ``subpipeline`` tasks (the child body) and ``@task`` tasks (the + traced signature), unioned with observed names defensively. Opaque + ``ref`` / ``@registered`` tasks get only their observed names (every + supplied argument, edge or literal, plus consumed outputs), marked + ``approximate``; nothing is read remotely. + """ + task_refs = dict(builder.task_refs_for_local_from_python) + interfaces: dict[str, TaskInterface] = {} + for task_id, (observed_inputs, observed_outputs) in _observed_task_ports(artifact.body).items(): + known: TaskInterface | None = None + child = artifact.subpipeline_children.get(task_id) + ref = task_refs.get(task_id) + if child is not None: + known = TaskInterface( + inputs=_graph_io_names(child.body, "inputs"), + outputs=_graph_io_names(child.body, "outputs"), + ) + elif ref is not None and sidecar_plan is not None: + fragment = sidecar_plan.fragment_by_task.get(task_id) + entry = sidecar_plan.entries.get(fragment, {}) if fragment is not None else {} + unwrapped = (entry.get("local_from_python") or {}).get("unwrapped_inputs") + known = _task_interface_for_ref(ref, unwrapped) + interfaces[task_id] = _union_interface(known, observed_inputs, observed_outputs) + return interfaces + + +def _layout_items(body: Any) -> list[tuple[str, Any]]: + """``(label, item)`` for every annotatable layout target in ``body``.""" + items: list[tuple[str, Any]] = [] + if not isinstance(body, dict): + return items + implementation = body.get("implementation") + graph = implementation.get("graph") if isinstance(implementation, dict) else None + tasks = graph.get("tasks") if isinstance(graph, dict) else None + if isinstance(tasks, dict): + items.extend((f"task {task_id!r}", task) for task_id, task in tasks.items()) + for key in ("inputs", "outputs"): + specs = body.get(key) + if isinstance(specs, list): + items.extend((f"{key[:-1]} #{index}", spec) for index, spec in enumerate(specs)) + return items + + +def _strip_positions(body: Any, reference: Any) -> Any: + """Copy of ``body`` without ``editor.position`` on layout targets. + + An ``annotations`` block left empty is dropped when the same item in + ``reference`` had none or only a position, so a transform may create or + remove a position-only block, but not drop an explicit empty one. + """ + stripped = copy.deepcopy(body) + reference_items = dict(_layout_items(reference)) + for label, item in _layout_items(stripped): + if not isinstance(item, dict): + continue + annotations = item.get("annotations") + if not isinstance(annotations, dict): + continue + annotations.pop(POSITION_ANNOTATION, None) + original = reference_items.get(label) + original_annotations = original.get("annotations") if isinstance(original, dict) else None + position_only = isinstance(original_annotations, dict) and set(original_annotations) == {POSITION_ANNOTATION} + if not annotations and (original_annotations is None or position_only): + del item["annotations"] + return stripped + + +def _apply_layout_transform( + artifact: SubgraphArtifact, + builder: Any, + sidecar_plan: TaskSidecarPlan | None, + ctx: CompileContext, + *, + layout: Layout, + occurrence_path: tuple[str, ...], +) -> None: + """Run the compile's layout transform on one artifact body, in place. + + Called once per artifact variant, after its refs are final and its + subgraphs are compiled, before validation/writing. The transform gets a + private deep copy and must return a dict that differs only in + ``editor.position`` annotations (string values) on graph tasks and + top-level inputs/outputs; anything else is a :class:`CompileError` naming + the location, and nothing is written. + """ + name = artifact.key.pipeline_name + where = f"pipeline {name!r} (occurrence path {occurrence_path!r})" + context = GraphLayoutContext( + layout=layout, + path=occurrence_path, + pipeline_name=name, + task_interfaces=types.MappingProxyType(_layout_task_interfaces(artifact, builder, sidecar_plan)), + artifact_dir=artifact.output_path.parent, + ) + before = artifact.body + try: + after = ctx.layout_transform(copy.deepcopy(before), context) + except CompileError: + raise + except Exception as exc: + raise CompileError(f"layout transform failed for {where}: {type(exc).__name__}: {exc}") from exc + if type(after) is not dict: + raise CompileError(f"layout transform for {where} must return a dict pipeline body.") + # Detach from anything the transform may have retained, so a later + # mutation of its own reference cannot reach the guarded, written body. + after = copy.deepcopy(after) + for label, item in _layout_items(after): + annotations = item.get("annotations") if isinstance(item, dict) else None + if isinstance(annotations, dict) and POSITION_ANNOTATION in annotations: + if type(annotations[POSITION_ANNOTATION]) is not str: + raise CompileError( + f"layout transform for {where} set a non-string {POSITION_ANNOTATION} on {label}." + ) + if _strip_positions(after, before) != _strip_positions(before, before): + raise CompileError( + f"layout transform for {where} changed the pipeline beyond " + f"{POSITION_ANNOTATION} annotations on graph tasks and inputs/outputs." + ) + artifact.body = after + + # --------------------------------------------------------------------------- # Recursive in-memory compile @@ -427,6 +750,8 @@ def _compile_pipeline_fn( base_dir: Path, rebroadcast_overrides: Mapping[str, Any] | None = None, precomputed_key: PipelineCompileKey | None = None, + layout_policy: _LayoutPolicy = _NO_LAYOUT, + occurrence_path: tuple[str, ...] = (), ) -> SubgraphArtifact: """Trace + emit ``pipeline_fn`` into an in-memory :class:`SubgraphArtifact`. @@ -446,6 +771,9 @@ def _compile_pipeline_fn( 4. Compile any ``subpipeline(child)`` children recursively and rewrite the parent task refs to pure ``file://`` URLs. + 5. When ``layout_policy`` covers this graph, run the layout transform + on the final body (post-order: subgraphs are already laid out). + Returns the planned artifact; validation + writing happen later in :func:`compile_pipeline` once the whole bundle is built. """ @@ -513,6 +841,7 @@ def _compile_pipeline_fn( # placeholder never reaches output. components_entries: dict[str, Any] = {} components_path: Path | None = None + sidecar_plan: TaskSidecarPlan | None = None task_refs = builder.task_refs_for_local_from_python if task_refs and not ctx.emit_components_sidecar: task_ids = sorted(tid for tid, _ref in task_refs) @@ -637,12 +966,30 @@ def _compile_pipeline_fn( ctx.broadcast_stack.append(BroadcastLayer(config={**raw_cfg, **dict(rebroadcast)})) pushed_broadcast = True try: - _process_subpipeline_children(artifact, builder, ctx, parent_propagate_config=pipeline_fn.propagate_config) + _process_subpipeline_children( + artifact, + builder, + ctx, + parent_propagate_config=pipeline_fn.propagate_config, + layout_carry=layout_policy.carry, + occurrence_path=occurrence_path, + ) finally: ctx.active_stack.pop() if pushed_broadcast: ctx.broadcast_stack.pop() + # 7. Compile-time layout, once per artifact variant, on the final body. + if layout_policy.apply is not None: + _apply_layout_transform( + artifact, + builder, + sidecar_plan, + ctx, + layout=layout_policy.apply, + occurrence_path=occurrence_path, + ) + return artifact @@ -652,6 +999,8 @@ def _process_subpipeline_children( ctx: CompileContext, *, parent_propagate_config: bool, + layout_carry: Layout | None = None, + occurrence_path: tuple[str, ...] = (), ) -> None: """Compile each ``subpipeline(child)(...)`` recorded during ``builder``'s trace and rewrite the parent task's ``componentRef`` to a pure @@ -772,7 +1121,13 @@ def _process_subpipeline_children( "published component boundary." ) - child_artifact = ctx.registry.get(child_key) + # Layout policy of THIS occurrence (the edge's own PipelineFn wins over + # the inherited one). It selects the artifact VARIANT; cycle detection + # above deliberately uses the layout-free ``child_key``. + child_policy = _layout_policy_for(child_fn, layout_carry, ctx) + child_variant = child_policy.variant() + registry_key = (child_key, child_variant) + child_artifact = ctx.registry.get(registry_key) if child_artifact is None: # Max-depth guard: the chain TO the child would be one deeper # than the current stack. @@ -785,7 +1140,7 @@ def _process_subpipeline_children( "pipeline graph." ) slug = _slugify(child_fn.name) - child_output_path = ctx.subgraph_dir / f"{slug}-{child_key.hash8()}.yaml" + child_output_path = ctx.subgraph_dir / f"{slug}-{_variant_hash8(child_key, child_variant)}.yaml" # Per-edge explicit-override depth: an # explicit ``.override_config`` set on THIS edge flows deep iff the # CALLER (this parent) is flagged. When it is, push a broadcast @@ -820,17 +1175,20 @@ def _process_subpipeline_children( base_dir=child_base_dir, rebroadcast_overrides={}, precomputed_key=child_key, + layout_policy=child_policy, + occurrence_path=(*occurrence_path, task_id), ) finally: if pushed_edge_override: ctx.broadcast_stack.pop() - ctx.registry[child_key] = child_artifact + ctx.registry[registry_key] = child_artifact parent_artifact.children.append(child_artifact) else: # Reused (dedup / diamond) — still a child of this parent for # structural completeness, but compiled only once. if child_artifact not in parent_artifact.children: parent_artifact.children.append(child_artifact) + parent_artifact.subpipeline_children[task_id] = child_artifact # Rewrite this task's componentRef to a pure file:// URL relative to # the REFERENCING artifact's directory. @@ -2866,6 +3224,7 @@ def compile_file( emit_components_sidecar: bool = True, image_overrides: Mapping[str, str] | None = None, pipeline_annotations: Mapping[str, str] | None = None, + layout_transform: GraphLayoutTransform | None = None, ) -> CompileResult: """Compile ``script`` to a single dehydrated pipeline YAML at ``output``. @@ -2886,6 +3245,7 @@ def compile_file( emit_components_sidecar=emit_components_sidecar, image_overrides=image_overrides, pipeline_annotations=pipeline_annotations, + layout_transform=layout_transform, ) 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 b5d8ba2..7b8a975 100644 --- a/packages/tangle-cli/src/tangle_cli/pipelines.py +++ b/packages/tangle-cli/src/tangle_cli/pipelines.py @@ -361,6 +361,7 @@ def compile_pipeline_file( image_overrides: Mapping[str, str] | None = None, pipeline_annotations: Mapping[str, str] | None = None, logger: Any | None = None, + layout_transform: Any | None = None, ) -> CompileResult: """Compile a Python-authored pipeline to a dehydrated YAML bundle. @@ -378,6 +379,10 @@ def compile_pipeline_file( 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`. + + ``layout_transform`` is the optional compile-time + :class:`~tangle_cli.python_pipeline.GraphLayoutTransform` for ``@Layout()`` + graphs, forwarded unchanged. """ from .pipeline_compiler import PipelineCompiler @@ -397,6 +402,7 @@ def compile_pipeline_file( # 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, + layout_transform=layout_transform, ) except (CompileError, SchemaValidationError) as exc: raise PipelineValidationError(str(exc)) from exc diff --git a/packages/tangle-cli/src/tangle_cli/python_pipeline/__init__.py b/packages/tangle-cli/src/tangle_cli/python_pipeline/__init__.py index b432643..a1e7525 100644 --- a/packages/tangle-cli/src/tangle_cli/python_pipeline/__init__.py +++ b/packages/tangle-cli/src/tangle_cli/python_pipeline/__init__.py @@ -4,6 +4,10 @@ from tangle_cli.python_pipeline import pipeline, task, registered, ref, raw, subpipeline, TaskEnv, In, Out +``Layout`` (``@Layout()`` / ``@Layout("name", recursive=...)``) requests +compile-time graph auto-layout, performed by a caller-supplied +``GraphLayoutTransform``; without one it does not change compiled YAML. + ``cfg`` is NOT a top-level export — it is a parameter the framework injects into the user's pipeline function at trace time. Importing the :class:`tangle_cli.python_pipeline.cfg.Cfg` class is reserved for the @@ -21,7 +25,9 @@ from __future__ import annotations from .dynamic_data import dynamic_secret +from .errors import InvalidLayoutError from .graph_io import graph_input, graph_output +from .layout import GraphLayoutContext, GraphLayoutTransform, Layout, TaskInterface from .pipeline import pipeline from .raw import raw from .ref import ref @@ -47,4 +53,9 @@ "In", "Out", "Outputs", + "Layout", + "GraphLayoutContext", + "GraphLayoutTransform", + "TaskInterface", + "InvalidLayoutError", ] 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 118620f..bc0589d 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 @@ -231,6 +231,9 @@ class SubgraphArtifact: # Child sub-artifacts this artifact references (for structure only — # the authoritative write/dedup list is ``CompileContext.registry``). children: list["SubgraphArtifact"] = field(default_factory=list) + # ``subpipeline(...)`` task ID -> the child artifact it runs, one entry + # per task (dedup hits included). Used for layout task interfaces. + subpipeline_children: dict[str, "SubgraphArtifact"] = field(default_factory=dict) # Cached dumped YAML text, filled during validation so the write pass # does not re-dump (and writes the exact bytes that were validated). dumped_text: str | None = None @@ -262,9 +265,12 @@ class CompileContext: # 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. - registry: dict[PipelineCompileKey, SubgraphArtifact] = field(default_factory=dict) + # Compiled CHILD artifacts keyed by ``(compile key, layout policy)`` + # (Decision M dedup). The policy is ``None`` unless a layout transform is + # installed and a ``@Layout()`` covers the child, so plain compiles dedup + # exactly as before; a child laid out under two different policies is + # two artifacts. The root is NOT stored here; it is returned directly. + registry: dict[tuple[PipelineCompileKey, str | None], SubgraphArtifact] = field(default_factory=dict) # Keys on the current recursive compile chain, for cycle detection # (Decision L). Pushed before recursing into a child, popped after. active_stack: list[PipelineCompileKey] = field(default_factory=list) @@ -284,3 +290,8 @@ class CompileContext: # not reuse a stale cached sibling module (the P2 sibling-import leak). source_dirs: set[Path] = field(default_factory=set) warnings: list[str] = field(default_factory=list) + # Compile-time ``GraphLayoutTransform`` (``None``: ``@Layout()`` has no + # effect on output and only produces a warning). + layout_transform: Any = None + # Display names of decorated pipelines already warned about (no transform). + layout_warned: set[str] = field(default_factory=set) 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 a4eed4a..e4c98a5 100644 --- a/packages/tangle-cli/src/tangle_cli/python_pipeline/errors.py +++ b/packages/tangle-cli/src/tangle_cli/python_pipeline/errors.py @@ -38,6 +38,16 @@ class InvalidEditorLayoutError(CompileError): """ +class InvalidLayoutError(CompileError): + """Raised on a misused ``@Layout()`` graph auto-layout request. + + Covers bare ``@Layout`` / positional arguments, an ``algorithm`` that is + not ``None`` or a non-empty ``str``, an unsupported decoration target, + and a second ``@Layout()`` on one graph definition. Messages never echo + the rejected value or target. + """ + + class InvalidInputDefaultError(CompileError): """Raised on an ``In[T]`` parameter default the schema cannot carry. diff --git a/packages/tangle-cli/src/tangle_cli/python_pipeline/layout.py b/packages/tangle-cli/src/tangle_cli/python_pipeline/layout.py new file mode 100644 index 0000000..b5b5a18 --- /dev/null +++ b/packages/tangle-cli/src/tangle_cli/python_pipeline/layout.py @@ -0,0 +1,237 @@ +"""``@Layout()`` — compile-time graph auto-layout for a pipeline definition. + +``Layout`` asks the compiler to auto-layout a graph, optionally with a named +algorithm:: + + from tangle_cli.python_pipeline import Layout, pipeline + + @Layout("banded") # this graph and its subgraphs + @pipeline("Daily Pulse") + def daily_pulse(...): ... + + @Layout(recursive=False) # this graph only + @pipeline("Judge") + def judge(...): ... + +Both stacking orders around ``@pipeline`` are supported; the decorator +returns its target unchanged. This package ships no layout algorithm: the +caller of :func:`tangle_cli.pipeline_compiler.compile_pipeline` passes a +:class:`GraphLayoutTransform`, which receives each covered graph body and +returns it with ``editor.position`` annotations. Without a transform the +request has no effect on output (the compiler emits a warning). + +Validation here is structural only: ``algorithm`` must be ``None`` or an +exact, non-empty ``str``, and ``recursive`` an exact ``bool``. Which names +are supported, and what ``None`` (the default) means, is the transform's +decision. Error messages never render the rejected value or target. +""" + +from __future__ import annotations + +import types +from collections.abc import Mapping +from dataclasses import dataclass +from pathlib import Path +from typing import Any, Protocol, TypeVar + +from .errors import InvalidLayoutError + +_T = TypeVar("_T") + +# Private marker for the ``@pipeline`` -over- ``@Layout()`` order: the inner +# ``Layout`` stamps the plain function, and ``@pipeline`` snapshots it onto +# the ``PipelineFn`` it builds (see :func:`layout_marker_of`). +_LAYOUT_MARKER = "__tangle_layout__" + +_BARE_OR_POSITIONAL = ( + "Layout must be called, with at most one positional algorithm name: use " + '@Layout(), @Layout("") or @Layout(algorithm="", recursive=...).' +) +_BAD_ALGORITHM = "Layout(algorithm=...) must be None or a non-empty str." +_BAD_RECURSIVE = "Layout(recursive=...) must be True or False." +_BAD_TARGET = ( + "@Layout() decorates a graph definition only: a @pipeline function, or a " + "plain function that @pipeline then wraps. Task, subpipeline and other " + "objects are not supported; decorate the child @pipeline definition instead." +) +_DUPLICATE = ( + "this pipeline definition already has a @Layout(); apply @Layout() at " + "most once per graph definition." +) + + +class Layout: + """Immutable request to auto-layout a graph definition at compile time. + + Args: + algorithm: Named layout algorithm, or ``None`` for the transform's + default. May be given positionally (an exact ``str`` only) or by + keyword. A bare ``@Layout`` without parentheses is refused. + recursive: When ``True`` (the default) the request also covers every + undecorated descendant subgraph. ``False`` covers only this + graph. A descendant with its own ``@Layout()`` always uses its + own declaration, and its own ``recursive`` governs below it. + """ + + __slots__ = ("_algorithm", "_recursive") + + _algorithm: str | None + _recursive: bool + + def __init__(self, *args: Any, algorithm: str | None = None, recursive: bool = True) -> None: + if args: + # A bare ``@Layout`` passes its target positionally; only one + # exact-``str`` positional (the algorithm) is accepted, and only + # without the keyword. The argument is never inspected or rendered. + if len(args) != 1 or type(args[0]) is not str or algorithm is not None: + raise InvalidLayoutError(_BARE_OR_POSITIONAL) + algorithm = args[0] + if algorithm is not None and (type(algorithm) is not str or not algorithm): + raise InvalidLayoutError(_BAD_ALGORITHM) + if type(recursive) is not bool: + raise InvalidLayoutError(_BAD_RECURSIVE) + object.__setattr__(self, "_algorithm", algorithm) + object.__setattr__(self, "_recursive", recursive) + + @property + def algorithm(self) -> str | None: + """The requested algorithm name, or ``None`` for the transform's default.""" + return self._algorithm + + @property + def recursive(self) -> bool: + """Whether undecorated descendant subgraphs are covered too.""" + return self._recursive + + def __setattr__(self, name: str, value: Any) -> None: + raise AttributeError("Layout is immutable") + + def __delattr__(self, name: str) -> None: + raise AttributeError("Layout is immutable") + + def __eq__(self, other: object) -> bool: + if type(other) is not Layout: + return NotImplemented + return (self._algorithm, self._recursive) == (other._algorithm, other._recursive) + + def __hash__(self) -> int: + return hash((Layout, self._algorithm, self._recursive)) + + def __repr__(self) -> str: + parts = [] + if self._algorithm is not None: + parts.append(f"algorithm={self._algorithm!r}") + if not self._recursive: + parts.append("recursive=False") + return f"Layout({', '.join(parts)})" + + def __reduce__(self) -> tuple[Any, ...]: + return (_rebuild_layout, (self._algorithm, self._recursive)) + + def __call__(self, target: _T) -> _T: + # Local import: ``pipeline`` imports this module at load time. + from .pipeline import PipelineFn + + # Exact types only, so no author-defined attribute or container hook + # runs: subclasses of PipelineFn are refused (as the deployment + # decorators do), and a function's ``__dict__`` is accessed through + # unbound ``dict`` methods in case it was replaced by a dict subclass. + target_type = type(target) + if target_type is PipelineFn: + if target.layout is not None: # type: ignore[attr-defined] + raise InvalidLayoutError(_DUPLICATE) + target.layout = self # type: ignore[attr-defined] + return target + if target_type is types.FunctionType: + namespace = target.__dict__ # type: ignore[attr-defined] + if dict.__contains__(namespace, _LAYOUT_MARKER): + raise InvalidLayoutError(_DUPLICATE) + dict.__setitem__(namespace, _LAYOUT_MARKER, self) + return target + raise InvalidLayoutError(_BAD_TARGET) + + +def _rebuild_layout(algorithm: str | None, recursive: bool = True) -> Layout: + return Layout(algorithm=algorithm, recursive=recursive) + + +def layout_marker_of(fn: Any) -> Layout | None: + """Return the ``@Layout()`` stamped directly on plain function ``fn``. + + Only exact Python functions are read, through their own ``__dict__`` via + unbound ``dict`` methods, so no author-defined attribute or container + hook runs. A marker that is not a genuine :class:`Layout` is ignored. + """ + if type(fn) is not types.FunctionType: + return None + marker = dict.get(fn.__dict__, _LAYOUT_MARKER) + return marker if type(marker) is Layout else None + + +@dataclass(frozen=True) +class TaskInterface: + """Ordered port names of one graph task, for layout geometry only. + + Attributes: + inputs: Input names: the component's declared inputs in declaration + order, followed by any other argument the task supplies. + outputs: Output names: declared outputs in declaration order, + followed by any other output consumed in the graph. + approximate: ``True`` when the component interface is not known + locally (``ref`` / ``@registered`` components): the inputs are + then every argument the task supplies (edges and filled literal + values) and the outputs those consumed elsewhere in the graph, so + the real card may be larger. ``False`` for ``@task`` and ``subpipeline`` tasks. + """ + + inputs: tuple[str, ...] + outputs: tuple[str, ...] + approximate: bool = False + + +@dataclass(frozen=True) +class GraphLayoutContext: + """What a :class:`GraphLayoutTransform` knows about the graph it lays out. + + Attributes: + layout: The governing :class:`Layout` (this graph's own declaration, + or the nearest recursive ancestor's). + path: Task-ID path from the root graph to the FIRST occurrence that + produced this artifact (``()`` for the root). Diagnostic only: + every occurrence with the same definition, config and layout + policy shares this artifact and this one transform call. + pipeline_name: The graph's ``@pipeline`` display name (diagnostic). + task_interfaces: Read-only map of EVERY graph task ID (in task order) + -> its :class:`TaskInterface`: exact for ``@task`` and + ``subpipeline`` tasks, observed-arguments/outputs-only (``approximate``) + for opaque ``ref`` / ``@registered`` tasks. Geometry input only: + interface data is never written back. + artifact_dir: Directory the artifact will be written to; relative + ``file://`` refs in the body resolve from here. The artifact + itself and its compiler-generated siblings are not written yet. + """ + + layout: Layout + path: tuple[str, ...] + pipeline_name: str + task_interfaces: Mapping[str, TaskInterface] + artifact_dir: Path + + @property + def algorithm(self) -> str | None: + """Shortcut for ``layout.algorithm``.""" + return self.layout.algorithm + + +class GraphLayoutTransform(Protocol): + """Compile-time layout callback passed to ``compile_pipeline``. + + Called once per covered artifact, with a private deep copy of the + dehydrated pipeline body, after its refs are final and its subgraphs are + compiled, and before validation and writing. It returns the body to + write, which may differ only in ``editor.position`` annotations on graph + tasks and on top-level ``inputs`` / ``outputs``. It must be deterministic: + every occurrence sharing the artifact reuses the result. + """ + + def __call__(self, graph: dict[str, Any], context: GraphLayoutContext) -> dict[str, Any]: ... 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 0f2803f..1e70eff 100644 --- a/packages/tangle-cli/src/tangle_cli/python_pipeline/pipeline.py +++ b/packages/tangle-cli/src/tangle_cli/python_pipeline/pipeline.py @@ -23,6 +23,7 @@ from . import emit from .errors import InvalidEditorLayoutError, InvalidPipelineLabelsError from .graph import GraphBuilder +from .layout import Layout, layout_marker_of @dataclass @@ -55,6 +56,11 @@ class PipelineFn: # Off by default (Decision F isolation is preserved). Metadata only — no # config is read at decoration time; the compile driver acts on it. propagate_config: bool = False + # Graph auto-layout INTENT from ``@Layout()`` (either stacking order). + # Authoring metadata only: never emitted into YAML and not part of the + # compile/dedup key. The compiler reports it per graph occurrence in + # ``CompileResult.layout_requests``; a consumer runs the layout. + layout: Layout | None = field(default=None, compare=False) # ------------------------------------------------------------------ # Calling the decorated PipelineFn directly is reserved for the @@ -213,6 +219,10 @@ def decorator(fn: Callable[..., Any]) -> PipelineFn: caller_dir=caller_dir, output_name=output_name, propagate_config=propagate_config, + # ``@pipeline`` over ``@Layout()``: snapshot the function's own + # marker into THIS wrapper. An outer ``@Layout()`` later sets + # only the wrapper it decorates. + layout=layout_marker_of(fn), ) return decorator diff --git a/pyproject.toml b/pyproject.toml index d457bfb..8a7eb43 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "tangle-cli" -version = "0.1.28" +version = "0.1.29" description = "CLI for Tangle, the open-source ML pipeline orchestration platform" readme = "README.md" authors = [ diff --git a/tests/test_layout_annotation.py b/tests/test_layout_annotation.py new file mode 100644 index 0000000..b9efac7 --- /dev/null +++ b/tests/test_layout_annotation.py @@ -0,0 +1,904 @@ +"""``@Layout()`` compile-time graph auto-layout. + +``Layout`` is immutable authoring metadata on a graph definition. The compiler +resolves which graph occurrences it covers (``recursive`` propagation, nearest +declaration wins) and, when a ``layout_transform`` is installed, lays each +covered artifact out before validation/writing. Without a transform, or +without ``@Layout()``, compiled bytes and filenames are unchanged. +""" + +from __future__ import annotations + +import copy +import functools +import json +import pickle +import sys +import textwrap +import types +from pathlib import Path +from typing import Any + +import pytest +import yaml + +import tangle_cli.python_pipeline as pp +from tangle_cli import pipeline_compiler +from tangle_cli.pipeline_compiler import PipelineCompiler, compile_pipeline +from tangle_cli.pipelines import PipelineValidationError, compile_pipeline_file +from tangle_cli.python_pipeline import ( + GraphLayoutContext, + InvalidLayoutError, + Layout, + TaskInterface, + pipeline, + subpipeline, + task, +) +from tangle_cli.python_pipeline.errors import CompileError +from tangle_cli.python_pipeline.pipeline import PipelineFn + +SECRET = "s3cr3t-value-never-rendered" +POSITION = "editor.position" + + +class _Hostile: + """Every rendering, comparison or attribute hook raises ``SECRET``.""" + + def _boom(self, *args: object, **kwargs: object) -> Any: # pragma: no cover - must never run + raise AssertionError(SECRET) + + __repr__ = __str__ = __format__ = __eq__ = __bool__ = __len__ = __getattr__ = __call__ = _boom + __hash__ = None # type: ignore[assignment] + + +class _HostileStr(str): + def __repr__(self) -> str: # pragma: no cover - must never run + raise AssertionError(SECRET) + + +def _fn(): + return None + + +@task(image="python:3.12") +def _some_task(x: str = "1"): + print(x) + + +class _Klass: + def method(self): + return None + + def __call__(self): + return None + + +class _TrappedPipelineFn(PipelineFn): + def __getattribute__(self, name: str) -> object: + if name == "layout": + raise RuntimeError(SECRET) + return super().__getattribute__(name) + + +class _TrappedDict(dict): + def _boom(self, *args: object) -> Any: # pragma: no cover - must never run + raise RuntimeError(SECRET) + + __contains__ = get = __getitem__ = __setitem__ = _boom + + +# --------------------------------------------------------------------------- +# The decorator + + +def test_layout_value_contract_and_exports(): + for name in ("Layout", "GraphLayoutContext", "GraphLayoutTransform", "TaskInterface", "InvalidLayoutError"): + assert name in pp.__all__ + assert getattr(pipeline_compiler, name, getattr(pp, name)) is getattr(pp, name) + assert issubclass(InvalidLayoutError, CompileError) + + assert (Layout().algorithm, Layout().recursive) == (None, True) + assert Layout("banded") == Layout(algorithm="banded") != Layout("banded", recursive=False) + assert repr(Layout()) == "Layout()" + assert repr(Layout("banded", recursive=False)) == "Layout(algorithm='banded', recursive=False)" + # Unsupported names are the transform's concern. + assert Layout("not-a-real-engine").algorithm == "not-a-real-engine" + + layout = Layout("banded", recursive=False) + for name in ("algorithm", "recursive", "_algorithm", "extra"): + with pytest.raises(AttributeError): + setattr(layout, name, "other") + with pytest.raises(AttributeError): + del layout._algorithm + assert hash(layout) == hash(Layout("banded", recursive=False)) + assert pickle.loads(pickle.dumps(layout)) == copy.deepcopy(layout) == layout + + +def test_layout_decorates_both_orders_and_snapshots_per_wrapper(): + # Outer order: the PipelineFn itself is returned, carries the layout, and + # keeps its other @pipeline metadata. + outer_target = pipeline("Outer", flow_direction="left-to-right", annotations={"a": "b"}, propagate_config=True)( + _fn + ) + banded = Layout("banded") + assert banded(outer_target) is outer_target and outer_target.layout is banded + assert outer_target.annotations == {"a": "b", "editor.flow-direction": "left-to-right"} + assert outer_target.propagate_config is True + + # Inner order: the plain function is returned; @pipeline snapshots it. + def inner(): + return None + + default = Layout() + assert default(inner) is inner + first, second = pipeline("First")(inner), pipeline("Second")(inner) + assert first.layout is second.layout is default + + # An outer decoration only affects the wrapper it decorates, not the + # function, and leaves PipelineFn equality alone. + def shared(): + return None + + plain, styled = pipeline("Same")(shared), pipeline("Same")(shared) + Layout("banded")(styled) + assert plain.layout is None and "__tangle_layout__" not in shared.__dict__ + assert plain == styled + + # A forged marker is ignored; a dict-subclass ``__dict__`` never dispatches. + def forged(): + return None + + forged.__dict__["__tangle_layout__"] = "banded" + assert pipeline("Forged")(forged).layout is None + + def trapped(): + return None + + trapped.__dict__ = _TrappedDict() + assert pipeline("Plain Trapped")(trapped).layout is None + assert banded(trapped) is trapped + assert pipeline("Inner Trapped")(trapped).layout is banded + + +_BAD_CONSTRUCTIONS = { + # bare ``@Layout`` (the target arrives positionally) and bad positionals + "bare_function": (lambda: Layout(_fn), "@Layout()"), + "bare_pipeline": (lambda: Layout(pipeline("Bare")(_fn)), "@Layout()"), + "bare_hostile": (lambda: Layout(_Hostile()), "@Layout()"), + "positional_none": (lambda: Layout(None), "@Layout()"), + "two_positionals": (lambda: Layout("a", "b"), "@Layout()"), + "positional_str_subclass": (lambda: Layout(_HostileStr("banded")), "@Layout()"), + "positional_and_keyword": (lambda: Layout("a", algorithm="b"), "@Layout()"), + # algorithm: exact non-empty str or None + "empty": (lambda: Layout(algorithm=""), "non-empty str"), + "int": (lambda: Layout(algorithm=1), "non-empty str"), + "bool": (lambda: Layout(algorithm=True), "non-empty str"), + "bytes": (lambda: Layout(algorithm=b"banded"), "non-empty str"), + "str_subclass": (lambda: Layout(algorithm=_HostileStr("banded")), "non-empty str"), + "hostile": (lambda: Layout(algorithm=_Hostile()), "non-empty str"), + "layout": (lambda: Layout(algorithm=Layout()), "non-empty str"), + # recursive: exact bool + "recursive_none": (lambda: Layout(recursive=None), "True or False"), + "recursive_int": (lambda: Layout(recursive=1), "True or False"), + "recursive_str": (lambda: Layout(recursive="false"), "True or False"), + "recursive_hostile": (lambda: Layout(recursive=_Hostile()), "True or False"), +} + + +def test_layout_refuses_invalid_construction_safely(): + for case, (build, fragment) in _BAD_CONSTRUCTIONS.items(): + with pytest.raises(InvalidLayoutError) as excinfo: + build() + message = str(excinfo.value) + assert fragment in message, case + assert SECRET not in message, case + with pytest.raises(TypeError): + Layout(algo="banded") # type: ignore[call-arg] + + +def test_layout_refuses_wrong_targets_and_duplicates_safely(): + child = pipeline("Child Target")(_fn) + wrong_targets = [ + _some_task, + subpipeline(child), + subpipeline(child).named("Handle"), + _Klass, + _Klass().method, + _Klass(), + functools.partial(_fn), + print, + None, + _Hostile(), + _TrappedPipelineFn(fn=_fn, name="Trapped"), + ] + for target in wrong_targets: + with pytest.raises(InvalidLayoutError) as excinfo: + Layout()(target) + assert "graph definition" in str(excinfo.value) + assert SECRET not in str(excinfo.value) + assert child.layout is None + + def fn(): + return None + + Layout()(fn) + pfn = Layout()(pipeline("Dup")(_fn)) + mixed = pipeline("Mixed")(fn) # inner marker snapshotted + trapped = _fn_with_trapped_dict() + Layout()(trapped) + for target in (fn, pfn, mixed, trapped): + with pytest.raises(InvalidLayoutError, match="at most once"): + Layout("banded")(target) + assert fn.__dict__["__tangle_layout__"] == pfn.layout == mixed.layout == Layout() + + +def _fn_with_trapped_dict(): + def fn(): + return None + + fn.__dict__ = _TrappedDict() + return fn + + +# --------------------------------------------------------------------------- +# Compile-time layout: helpers + +# Deterministic test transform: every task/input gets a position whose ``y`` +# encodes the algorithm, so each artifact's effective policy is observable in +# its written YAML. +_ALGO_Y = {None: 0, "sugiyama": 1000, "banded": 2000, "a": 3000, "b": 4000} + + +class RecordingTransform: + def __init__(self) -> None: + self.calls: list[GraphLayoutContext] = [] + + def __call__(self, graph: dict[str, Any], context: GraphLayoutContext) -> dict[str, Any]: + self.calls.append(context) + y = _ALGO_Y[context.algorithm] + for index, task_spec in enumerate(graph["implementation"]["graph"]["tasks"].values()): + task_spec.setdefault("annotations", {})[POSITION] = json.dumps({"x": index * 100, "y": y}) + for index, spec in enumerate(graph.get("inputs") or []): + spec.setdefault("annotations", {})[POSITION] = json.dumps({"x": -100, "y": y + index}) + return graph + + def summary(self) -> list[tuple[str, tuple[str, ...], str | None]]: + return [(c.pipeline_name, c.path, c.algorithm) for c in self.calls] + + +# Placeholders become ``Layout(...)`` or the line-count-preserving ``_same`` +# no-op, so decorated and plain sources differ only in decorators. +_BUNDLE = ''' +from tangle_cli.python_pipeline import Layout, In, Out, pipeline, subpipeline, task +import layout_probe + +layout_probe.events.append("root-module") + + +def _same(fn): + return fn + + +@task(image="python:3.12") +def leaf(value: str = "x"): + print(value) + + +{JUDGE} +@pipeline("Judge") +def judge(seed: In[str]) -> Out[str]: + layout_probe.events.append("trace-judge") + return leaf.named("Leaf")(value=seed) + + +@pipeline("Mid") +{MID} +def mid(seed: In[str]) -> Out[str]: + layout_probe.events.append("trace-mid") + first = subpipeline(judge).named("First")(seed=seed) + return subpipeline(judge).named("Second")(seed=first) + + +{ROOT} +@pipeline("Root") +def root(seed: In[str]) -> Out[str]: + layout_probe.events.append("trace-root") + middle = subpipeline(mid).named("Middle")(seed=seed) + return subpipeline(judge).named("Direct")(seed=middle) + + +@Layout(algorithm="unselected") +@pipeline("Unselected") +def unselected(seed: In[str]) -> Out[str]: + return subpipeline(judge).named("Never")(seed=seed) +''' + +_PATHS: list[tuple[str, ...]] = [(), ("Middle",), ("Middle", "First"), ("Middle", "Second"), ("Direct",)] + + +@pytest.fixture +def probe(monkeypatch): + module = types.ModuleType("layout_probe") + module.events = [] # type: ignore[attr-defined] + monkeypatch.setitem(sys.modules, "layout_probe", module) + return module + + +def _write(path: Path, source: str) -> Path: + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text(textwrap.dedent(source), encoding="utf-8") + return path + + +def _compile_bundle(tmp_path: Path, name: str, transform: Any = None, select: str = "root", **decorators: str): + values = {"JUDGE": "@_same", "MID": "@_same", "ROOT": "@_same", **decorators} + src = _write(tmp_path / "proj" / "bundle.py", _BUNDLE.format(**values)) + out = tmp_path / name / "compiled.yaml" + return compile_pipeline(src, out, pipeline_name=select, layout_transform=transform), out + + +def _bundle_bytes(out: Path) -> dict[str, bytes]: + return {p.relative_to(out.parent).as_posix(): p.read_bytes() for p in sorted(out.parent.rglob("*")) if p.is_file()} + + +def _load(path: Path) -> dict[str, Any]: + return yaml.safe_load(path.read_text(encoding="utf-8")) + + +def _occurrence(out: Path, *path: str) -> tuple[Path, dict[str, Any]]: + """Follow the written ``file://`` refs from the root to occurrence ``path``.""" + current = out.resolve() + body = _load(current) + for task_id in path: + url = body["implementation"]["graph"]["tasks"][task_id]["componentRef"]["url"] + current = (current.parent / url.removeprefix("file://")).resolve() + body = _load(current) + return current, body + + +def _ys(out: Path, *path: str) -> set[int | None]: + """Distinct laid-out ``y`` values of the tasks in the graph at ``path``.""" + _file, body = _occurrence(out, *path) + tasks = body["implementation"]["graph"]["tasks"].values() + values = [(t.get("annotations") or {}).get(POSITION) for t in tasks] + return {None if v is None else json.loads(v)["y"] for v in values} + + +# --------------------------------------------------------------------------- +# Coverage policy over one bundle: Root -> Middle(Mid) -> First/Second(Judge), +# Root -> Direct(Judge). Each row: decorators, expected y per occurrence path +# (None = not laid out), expected transform calls (post-order, once per +# artifact variant), sidecar count, and the paths whose sidecar must keep its +# legacy (no-layout) filename. + +_SUG, _BAN = '@Layout("sugiyama")', '@Layout("banded")' +_POLICY_CASES = { + "default_recursive_covers_subtree": ( + {"ROOT": _SUG}, + [1000, 1000, 1000, 1000, 1000], + [("Judge", ("Middle", "First"), "sugiyama"), ("Mid", ("Middle",), "sugiyama"), ("Root", (), "sugiyama")], + 2, + [], + ), + "recursive_false_covers_only_itself": ( + {"ROOT": '@Layout("sugiyama", recursive=False)'}, + [1000, None, None, None, None], + [("Root", (), "sugiyama")], + 2, + _PATHS[1:], + ), + "explicit_child_under_non_recursive_parent": ( + {"ROOT": "@Layout(recursive=False)", "JUDGE": _BAN}, + [0, None, 2000, 2000, 2000], + [("Judge", ("Middle", "First"), "banded"), ("Root", (), None)], + 2, + [("Middle",)], + ), + "recursive_override_splits_shared_child": ( + {"ROOT": _SUG, "MID": _BAN}, + [1000, 2000, 2000, 2000, 1000], + [ + ("Judge", ("Middle", "First"), "banded"), + ("Mid", ("Middle",), "banded"), + ("Judge", ("Direct",), "sugiyama"), + ("Root", (), "sugiyama"), + ], + 3, + [], + ), + "reset_to_default_is_inherited_below": ( + {"ROOT": _SUG, "MID": "@Layout()"}, + [1000, 0, 0, 0, 1000], + [ + ("Judge", ("Middle", "First"), None), + ("Mid", ("Middle",), None), + ("Judge", ("Direct",), "sugiyama"), + ("Root", (), "sugiyama"), + ], + 3, + [], + ), + "non_recursive_override_stops_all_inheritance": ( + {"ROOT": _SUG, "MID": '@Layout("banded", recursive=False)'}, + [1000, 2000, None, None, 1000], + [("Mid", ("Middle",), "banded"), ("Judge", ("Direct",), "sugiyama"), ("Root", (), "sugiyama")], + 3, + [("Middle", "First"), ("Middle", "Second")], + ), + "undecorated_root_with_decorated_descendant": ( + {"JUDGE": _BAN}, + [None, None, 2000, 2000, 2000], + [("Judge", ("Middle", "First"), "banded")], + 2, + [("Middle",)], + ), +} + + +@pytest.mark.parametrize("case", list(_POLICY_CASES)) +def test_compile_layout_coverage_policy(tmp_path, probe, case): + decorators, expected_ys, expected_calls, sidecars, legacy_paths = _POLICY_CASES[case] + _plain, plain_out = _compile_bundle(tmp_path, "plain") + del probe.events[:] + transform = RecordingTransform() + result, out = _compile_bundle(tmp_path, "out", transform, **decorators) + + assert transform.summary() == expected_calls + assert len(result.subgraph_paths) == sidecars + for path, expected in zip(_PATHS, expected_ys): + assert _ys(out, *path) == {expected}, path + # Root and subgraph graph inputs are positioned with their graph's tasks. + for path, expected in zip(_PATHS, expected_ys): + inputs = _occurrence(out, *path)[1]["inputs"] + position = (inputs[0].get("annotations") or {}).get(POSITION) + assert (position is None) == (expected is None), path + # Repeated calls under one policy share one artifact; uncovered children + # keep their legacy filenames. + assert _occurrence(out, "Middle", "First")[0] == _occurrence(out, "Middle", "Second")[0] + for path in legacy_paths: + assert _occurrence(out, *path)[0].name == _occurrence(plain_out, *path)[0].name, path + # One module execution; each definition is traced once per artifact + # variant (Judge has one or two), never once per occurrence. + assert [e for e in probe.events if e != "trace-judge"] == ["root-module", "trace-root", "trace-mid"] + judge_files = {_occurrence(out, *path)[0] for path in _PATHS[2:]} + assert probe.events.count("trace-judge") == len(judge_files) + + +def test_compile_layout_leaves_output_identical_without_transform_or_layout(tmp_path, probe): + _plain, plain_out = _compile_bundle(tmp_path, "plain") + plain_bytes = _bundle_bytes(plain_out) + + decorators = {"ROOT": _SUG, "MID": '@Layout("banded", recursive=False)', "JUDGE": "@Layout()"} + no_transform, out = _compile_bundle(tmp_path, "no-transform", None, **decorators) + assert _bundle_bytes(out) == plain_bytes + warned = [w for w in no_transform.warnings if "@Layout()" in w] + assert len(warned) == 3 and all(any(repr(n) in w for w in warned) for n in ("Root", "Mid", "Judge")) + + unused = RecordingTransform() + _result, out = _compile_bundle(tmp_path, "undecorated", unused) + assert unused.calls == [] and _bundle_bytes(out) == plain_bytes + + unselected = RecordingTransform() + _compile_bundle(tmp_path, "unselected", unselected, select="mid") + assert unselected.calls == [] + + # Laid-out output is deterministic across compiles. + _r1, out1 = _compile_bundle(tmp_path, "first", RecordingTransform(), ROOT=_SUG, MID=_BAN) + _r2, out2 = _compile_bundle(tmp_path, "second", RecordingTransform(), ROOT=_SUG, MID=_BAN) + assert _bundle_bytes(out1) == _bundle_bytes(out2) != plain_bytes + + +_DIAMOND = ''' +from tangle_cli.python_pipeline import Layout, In, Out, pipeline, subpipeline, task + +@task(image="python:3.12") +def leaf(value: str = "x"): + print(value) + +@pipeline("Bottom") +def bottom(seed: In[str]) -> Out[str]: + return leaf.named("Leaf")(value=seed) + +@Layout("{left}") +@pipeline("Left") +def left(seed: In[str]) -> Out[str]: + return subpipeline(bottom).named("Shared")(seed=seed) + +@Layout("{right}") +@pipeline("Right") +def right(seed: In[str]) -> Out[str]: + return subpipeline(bottom).named("Shared")(seed=seed) + +@pipeline("Top") +def top(seed: In[str]) -> Out[str]: + a = subpipeline(left).named("A")(seed=seed) + return subpipeline(right).named("B")(seed=a) +''' + + +@pytest.mark.parametrize(("left", "right", "bottom_calls", "sidecars"), [("a", "a", 1, 3), ("a", "b", 2, 4)]) +def test_compile_layout_diamond_dedups_same_policy_and_splits_different(tmp_path, left, right, bottom_calls, sidecars): + src = _write(tmp_path / "proj" / "diamond.py", _DIAMOND.format(left=left, right=right)) + transform = RecordingTransform() + out = tmp_path / "out" / "compiled.yaml" + result = compile_pipeline(src, out, pipeline_name="top", layout_transform=transform) + + assert [c.pipeline_name for c in transform.calls].count("Bottom") == bottom_calls + assert len(result.subgraph_paths) == sidecars + same_file = _occurrence(out, "A", "Shared")[0] == _occurrence(out, "B", "Shared")[0] + assert same_file == (left == right) + assert _ys(out, "A", "Shared") == {_ALGO_Y[left]} + assert _ys(out, "B", "Shared") == {_ALGO_Y[right]} + assert _ys(out) == {None} + + +def test_compile_layout_config_variants_and_edge_wrappers(tmp_path): + project = tmp_path / "proj" + _write(project / "child_config.yaml", "message: from-file\n") + src = _write( + project / "variants.py", + ''' + from tangle_cli.python_pipeline import Layout, In, Out, pipeline, subpipeline, task + + @task(image="python:3.12") + def emit(message: str = "x"): + print(message) + + @Layout("banded") + @pipeline("Variant Child", config="child_config.yaml") + def child(seed: In[str], cfg) -> Out[str]: + return emit.named("Emit")(message=cfg.message, wait_for=seed) + + def _body(seed: In[str]) -> Out[str]: + return emit.named("Emit")(message="plain", wait_for=seed) + + # Two wrappers of one function share a compile key; only one is styled. + plain_child = pipeline("Shared Child")(_body) + styled_child = Layout("banded")(pipeline("Shared Child")(_body)) + + @pipeline("Variant Parent") + def parent(seed: In[str]) -> Out[str]: + one = subpipeline(child).named("One").override_config(message="one")(seed=seed) + two = subpipeline(child).named("Two").override_config(message="two")(seed=one) + again = subpipeline(child).named("Again").override_config(message="one")(seed=two) + plain = subpipeline(plain_child).named("Plain")(seed=again) + return subpipeline(styled_child).named("Styled")(seed=plain) + ''', + ) + plain_out = tmp_path / "plain" / "compiled.yaml" + compile_pipeline(src, plain_out, pipeline_name="parent") + transform = RecordingTransform() + out = tmp_path / "out" / "compiled.yaml" + result = compile_pipeline(src, out, pipeline_name="parent", layout_transform=transform) + + assert transform.summary() == [ + ("Variant Child", ("One",), "banded"), + ("Variant Child", ("Two",), "banded"), + ("Shared Child", ("Styled",), "banded"), + ] + # Config variants stay distinct and the repeat dedups; the two wrappers + # split into an uncovered (legacy-named) and a laid-out artifact. + assert len(result.subgraph_paths) == 4 + assert _occurrence(out, "One")[0] == _occurrence(out, "Again")[0] != _occurrence(out, "Two")[0] + for path in [("One",), ("Two",), ("Again",), ("Styled",)]: + assert _ys(out, *path) == {2000}, path + assert _ys(out, "Plain") == {None} + assert _occurrence(out, "Plain")[0].name == _occurrence(plain_out, "Plain")[0].name + + +def test_compile_layout_cycle_detection_ignores_layout_variants(tmp_path): + """``Loop`` reached under ``Step``'s inherited layout is a different artifact + variant than the root ``Loop``, but still the same definition: the cycle + is reported at once, before any transform runs or anything is written.""" + src = _write( + tmp_path / "proj" / "cycle.py", + ''' + from tangle_cli.python_pipeline import Layout, In, Out, pipeline, subpipeline + + @pipeline("Loop") + def loop(seed: In[str]) -> Out[str]: + return subpipeline(step).named("Step")(seed=seed) + + @Layout("b") + @pipeline("Step") + def step(seed: In[str]) -> Out[str]: + return subpipeline(loop).named("Again")(seed=seed) + ''', + ) + transform = RecordingTransform() + out = tmp_path / "out" / "compiled.yaml" + with pytest.raises(CompileError) as excinfo: + compile_pipeline(src, out, pipeline_name="loop", layout_transform=transform) + message = str(excinfo.value) + assert "nested pipeline cycle detected" in message and "max depth" not in message + assert message.count("Loop (") == 2 and message.count("Step (") == 1 + assert transform.calls == [] and not out.parent.exists() + + +_IMPORTED_CHILD = ''' +from tangle_cli.python_pipeline import Layout, In, Out, pipeline, task +import layout_probe + +layout_probe.events.append("child-module") + + +@task(image="python:3.12") +def leaf(value: str = "x"): + print(value) + + +@pipeline("Imported Child") +@Layout(algorithm="banded") +def imported_child(seed: In[str]) -> Out[str]: + return leaf.named("Leaf")(value=seed) +''' + +_IMPORTING_ROOT = ''' +from tangle_cli.python_pipeline import In, Out, pipeline, subpipeline +from {module} import imported_child + + +@pipeline("Importing Root") +def importing_root(seed: In[str]) -> Out[str]: + return subpipeline(imported_child).named("Imported")(seed=seed) +''' + + +def test_compile_layout_imported_and_pre_imported_children(tmp_path, probe, monkeypatch): + # A sibling module imported by the root during compile. + _write(tmp_path / "proj" / "sibling_child.py", _IMPORTED_CHILD) + src = _write(tmp_path / "proj" / "root.py", _IMPORTING_ROOT.format(module="sibling_child")) + transform = RecordingTransform() + out = tmp_path / "sibling" / "compiled.yaml" + compile_pipeline(src, out, layout_transform=transform) + assert transform.summary() == [("Imported Child", ("Imported",), "banded")] + assert _ys(out, "Imported") == {2000} + + # A library module imported (and cached) BEFORE compile keeps its layout, + # across repeated compiles, without being re-executed. + module_name = "layout_cached_child_lib" + _write(tmp_path / "libs" / f"{module_name}.py", _IMPORTED_CHILD) + monkeypatch.syspath_prepend(str(tmp_path / "libs")) + monkeypatch.delitem(sys.modules, module_name, raising=False) + cached = __import__(module_name) + monkeypatch.setitem(sys.modules, module_name, cached) + src = _write(tmp_path / "cached-proj" / "root.py", _IMPORTING_ROOT.format(module=module_name)) + executions = probe.events.count("child-module") + for attempt in range(2): + sys.modules[module_name] = cached # the compiler purges it after each compile + transform = RecordingTransform() + out = tmp_path / f"cached{attempt}" / "compiled.yaml" + compile_pipeline(src, out, layout_transform=transform) + assert transform.summary() == [("Imported Child", ("Imported",), "banded")] + assert _ys(out, "Imported") == {2000} + assert probe.events.count("child-module") == executions + + +def test_compile_layout_context_and_task_interfaces(tmp_path): + project = tmp_path / "proj" + for name, spec in { + # Opaque ref components: their YAML is never read by the compiler, so + # declared-but-unused ports (``unused``, ``log``, ``hidden``) are absent. + "producer": "inputs: [{name: unused}, {name: seed}, {name: mode}]\noutputs: [{name: data}, {name: log}]", + "consumer": "inputs: [{name: mode}, {name: hidden}]", + "noop": "inputs: [{name: wait_for}]", + }.items(): + _write(project / f"{name}.yaml", f"name: {name}\n{spec}\nimplementation:\n container:\n image: alpine\n") + src = _write( + project / "interfaces.py", + ''' + from typing import NamedTuple + from tangle_cli.python_pipeline import Layout, In, Out, pipeline, ref, subpipeline, task + + class Pair(NamedTuple): + left: str + right: str + + @task(image="python:3.12") + def split(text: str, sep: str = ",") -> Pair: + return Pair(*text.split(sep, 1)) + + @pipeline("Inner") + def inner(seed: In[str], extra: In[str] = "e") -> Out[str]: + return split.named("Split")(text=seed) + + @Layout("banded") + @pipeline("Outer") + def outer(seed: In[str]) -> Out[str]: + produced = ref(url="file://./producer.yaml").named("Producer")(seed=seed, mode="fast") + child = subpipeline(inner).named("Child")(seed=produced.data) + ref(url="file://./consumer.yaml").named("Consumer")(mode="x", is_enabled=produced.flag) + ref(url="file://./noop.yaml").named("Noop")(wait_for=produced.data) + return child + ''', + ) + transform = RecordingTransform() + out = project / "compiled.yaml" # colocated so the author refs resolve + compile_pipeline(src, out, pipeline_name="outer", layout_transform=transform) + + inner_ctx, outer_ctx = transform.calls + assert (inner_ctx.pipeline_name, inner_ctx.path, outer_ctx.path) == ("Inner", ("Child",), ()) + assert inner_ctx.layout is outer_ctx.layout == Layout("banded") + assert inner_ctx.artifact_dir == (project / "compiled.subgraphs").resolve() + assert outer_ctx.artifact_dir == project.resolve() + assert dict(inner_ctx.task_interfaces) == { + # Exact signature ports, unioned with the observed bare-return output. + "Split": TaskInterface(inputs=("text", "sep"), outputs=("left", "right", "wait_for_output")), + } + assert list(outer_ctx.task_interfaces.items()) == [ + # Opaque: every supplied argument (edge AND literal) + consumed outputs, + # including one consumed only by another task's isEnabled condition. + ("Producer", TaskInterface(inputs=("seed", "mode"), outputs=("data", "flag"), approximate=True)), + ("Child", TaskInterface(inputs=("seed", "extra"), outputs=("wait_for_output",))), + ("Consumer", TaskInterface(inputs=("mode",), outputs=(), approximate=True)), + ("Noop", TaskInterface(inputs=("wait_for",), outputs=(), approximate=True)), + ] + with pytest.raises(TypeError): + outer_ctx.task_interfaces["Noop"] = TaskInterface((), ()) # type: ignore[index] + # Extra observed INPUT names cannot be authored for known interfaces (the + # compiler validates them), so the defensive input union is checked directly. + known = TaskInterface(inputs=("a", "b"), outputs=("x",)) + assert pipeline_compiler._union_interface(known, ["b", "c"], []) == TaskInterface(("a", "b", "c"), ("x",)) + # Opaque tasks are positioned too; no interface metadata is written. + assert _ys(out) == {2000} + assert "approximate" not in out.read_text(encoding="utf-8") + + +def test_compile_layout_transform_is_forwarded_by_every_entry_point(tmp_path, probe): + outs = [] + for name, run in { + "function": compile_pipeline, + "handler": PipelineCompiler().compile_file, + "pipelines_api": compile_pipeline_file, + }.items(): + values = {"JUDGE": _BAN, "MID": "@_same", "ROOT": "@_same"} + src = _write(tmp_path / "proj" / "bundle.py", _BUNDLE.format(**values)) + transform = RecordingTransform() + out = tmp_path / name / "compiled.yaml" + run(src, out, pipeline_name="root", layout_transform=transform) + assert len(transform.calls) == 1, name + outs.append(_bundle_bytes(out)) + assert outs[0] == outs[1] == outs[2] + + +# --------------------------------------------------------------------------- +# Transform contract guard + +_GUARDED = ''' +from tangle_cli.python_pipeline import Layout, In, Out, pipeline, task + +@task(image="python:3.12") +def leaf(value: str = "x"): + print(value) + +@Layout() +@pipeline("Single") +def single(seed: In[str]) -> Out[str]: + first = leaf.named("First").with_position(5, 5, width=10, height=10)(value=seed) + second = leaf.named("Second").with_annotations({"keep": "me"})(value=seed, wait_for=first) + return leaf.named("Third")(value=seed, wait_for=second) +''' + + +def _tasks(graph: dict[str, Any]) -> dict[str, Any]: + return graph["implementation"]["graph"]["tasks"] + + +def _set_positions(graph, context): + for index, task_spec in enumerate(_tasks(graph).values()): + task_spec.setdefault("annotations", {})[POSITION] = json.dumps({"x": index, "y": 7}) + return graph + + +def _remove_position_only_block(graph, context): + del _tasks(graph)["First"]["annotations"] + return graph + + +def _raise(exc: Exception): + def transform(graph, context): + raise exc + + return transform + + +class _Refused(CompileError): + pass + + +def _mutate(edit): + def transform(graph, context): + edit(graph) + return graph + + return transform + + +_GUARD_CASES = { + # allowed: create, replace (dropping manual width/height) and remove positions + "set_positions": (_set_positions, None), + "remove_position_only_block": (_remove_position_only_block, None), + # refused: anything beyond editor.position, bad values and bad returns + "rename": (_mutate(lambda g: g.update(name="Renamed")), "beyond"), + "argument": (_mutate(lambda g: _tasks(g)["Second"]["arguments"].update(value="changed")), "beyond"), + "other_annotation": (_mutate(lambda g: _tasks(g)["Third"].setdefault("annotations", {}).update(o="x")), "beyond"), + "drop_other_annotations": (_mutate(lambda g: _tasks(g)["Second"].pop("annotations")), "beyond"), + "drop_task": (_mutate(lambda g: _tasks(g).pop("Third")), "beyond"), + "non_string_position": ( + _mutate(lambda g: _tasks(g)["Third"].setdefault("annotations", {}).update({POSITION: {"x": 1}})), + "non-string", + ), + "return_none": (lambda graph, context: None, "must return a dict"), + "return_mapping_proxy": (lambda graph, context: types.MappingProxyType(graph), "must return a dict"), + # transform failures: CompileError propagates as-is, others are wrapped + "raises_value_error": (_raise(ValueError("unsupported algorithm")), "ValueError: unsupported algorithm"), + "raises_compile_error": (_raise(_Refused("flow direction not supported")), "flow direction not supported"), +} + + +def test_compile_layout_transform_guard(tmp_path): + src = _write(tmp_path / "proj" / "single.py", _GUARDED) + for case, (transform, error) in _GUARD_CASES.items(): + out = tmp_path / case / "compiled.yaml" + if error is None: + compile_pipeline(src, out, layout_transform=transform) + continue + with pytest.raises(CompileError) as excinfo: + compile_pipeline(src, out, layout_transform=transform) + assert error in str(excinfo.value), case + assert "'Single'" in str(excinfo.value) or case == "raises_compile_error", case + assert not out.parent.exists(), case + assert isinstance(excinfo.value, _Refused) + + tasks = _tasks(_load(tmp_path / "set_positions" / "compiled.yaml")) + assert [json.loads(t["annotations"][POSITION]) for t in tasks.values()] == [ + {"x": 0, "y": 7}, + {"x": 1, "y": 7}, + {"x": 2, "y": 7}, + ] + assert tasks["Second"]["annotations"]["keep"] == "me" + removed = _tasks(_load(tmp_path / "remove_position_only_block" / "compiled.yaml")) + assert "annotations" not in removed["First"] and removed["Second"]["annotations"] == {"keep": "me"} + + with pytest.raises(PipelineValidationError): + compile_pipeline_file(src, tmp_path / "api" / "compiled.yaml", layout_transform=_raise(ValueError("x"))) + + +def test_compile_layout_retained_transform_references_cannot_reach_output(tmp_path, probe): + """A transform that keeps the bodies it was given/returned cannot change + the guarded, written output later (from the parent's call or afterwards).""" + retained: list[dict[str, Any]] = [] + + def sneaky(graph, context): + for earlier in retained: + for task_spec in _tasks(earlier).values(): + task_spec.setdefault("arguments", {})["value"] = "CHANGED_AFTER_GUARD" + result = RecordingTransform()(graph, context) + retained.append(result) + return result + + _result, out = _compile_bundle(tmp_path, "out", sneaky, ROOT=_SUG) + assert len(retained) == 3 + retained[-1]["name"] = "CHANGED_AFTER_COMPILE" + for text in (b.decode("utf-8") for b in _bundle_bytes(out).values()): + assert "CHANGED_AFTER" not in text + + +def test_compile_layout_misuse_in_module_fails_before_writing(tmp_path): + src = _write( + tmp_path / "dup.py", + ''' + from tangle_cli.python_pipeline import Layout, Out, pipeline + + @Layout() + @pipeline("Dup Root") + @Layout() + def dup() -> Out[str]: + return None + ''', + ) + out = tmp_path / "out" / "compiled.yaml" + with pytest.raises(CompileError, match="at most once"): + compile_pipeline(src, out, layout_transform=RecordingTransform()) + assert not out.parent.exists() diff --git a/tests/test_packaging.py b/tests/test_packaging.py index fb1333a..3f41ed4 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.28" in metadata + assert "Version: 0.1.29" 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_python_pipeline_dsl.py b/tests/test_python_pipeline_dsl.py index d2e826c..c52f911 100644 --- a/tests/test_python_pipeline_dsl.py +++ b/tests/test_python_pipeline_dsl.py @@ -65,6 +65,11 @@ def test_all_names_are_exported(self): "In", "Out", "Outputs", + "Layout", + "GraphLayoutContext", + "GraphLayoutTransform", + "TaskInterface", + "InvalidLayoutError", } def test_every_all_name_is_present_on_module(self): diff --git a/uv.lock b/uv.lock index 7d562b1..296015c 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.28" +version = "0.1.29" source = { editable = "." } dependencies = [ { name = "cloud-pipelines" },