kaxil commented on code in PR #74084:
URL: https://github.com/apache/airflow/pull/74084#discussion_r4183857427


##########
task-sdk/src/airflow/sdk/importers/yaml_importer/schema.json:
##########
@@ -0,0 +1,306 @@
+{
+  "$defs": {
+    "CodeTask": {
+      "additionalProperties": true,
+      "description": "A task running custom code.\n\nThis should have ``run`` 
holding function arguments, and optionally\n``queue`` to route it.",
+      "properties": {
+        "id": {
+          "title": "Id",
+          "type": "string"
+        },
+        "needs": {
+          "items": {
+            "type": "string"
+          },
+          "title": "Needs",
+          "type": "array"
+        },
+        "extends": {
+          "items": {
+            "type": "string"
+          },
+          "title": "Extends",
+          "type": "array"
+        },
+        "with": {
+          "additionalProperties": {
+            "oneOf": [
+              {
+                "$ref": "#/$defs/XComRef"
+              },
+              {
+                "$ref": "#/$defs/TemplateRef"
+              },
+              {
+                "$ref": "#/$defs/ConstRef"
+              },
+              {
+                "$ref": "#/$defs/_Literal"
+              }
+            ]
+          },
+          "title": "With",
+          "type": "object"
+        },
+        "run": {
+          "additionalProperties": {
+            "oneOf": [
+              {
+                "$ref": "#/$defs/XComRef"
+              },
+              {
+                "$ref": "#/$defs/TemplateRef"
+              },
+              {
+                "$ref": "#/$defs/ConstRef"
+              },
+              {
+                "$ref": "#/$defs/_Literal"
+              }
+            ]
+          },
+          "title": "Run",
+          "type": "object"
+        }
+      },
+      "required": [
+        "id"
+      ],
+      "title": "CodeTask",
+      "type": "object"
+    },
+    "ConstRef": {
+      "additionalProperties": false,
+      "description": "Force a literal; the value is taken verbatim.",
+      "properties": {
+        "$const": {
+          "title": "$Const"
+        }
+      },
+      "required": [
+        "$const"
+      ],
+      "title": "ConstRef",
+      "type": "object"
+    },
+    "OperatorTask": {
+      "additionalProperties": true,
+      "description": "A ready-made operator.\n\nThis should have ``uses`` 
(import path) + ``with``.",
+      "properties": {
+        "id": {
+          "title": "Id",
+          "type": "string"
+        },
+        "needs": {
+          "items": {
+            "type": "string"
+          },
+          "title": "Needs",
+          "type": "array"
+        },
+        "extends": {
+          "items": {
+            "type": "string"
+          },
+          "title": "Extends",
+          "type": "array"
+        },
+        "with": {
+          "additionalProperties": {
+            "oneOf": [
+              {
+                "$ref": "#/$defs/XComRef"
+              },
+              {
+                "$ref": "#/$defs/TemplateRef"
+              },
+              {
+                "$ref": "#/$defs/ConstRef"
+              },
+              {
+                "$ref": "#/$defs/_Literal"
+              }
+            ]
+          },
+          "title": "With",
+          "type": "object"
+        },
+        "uses": {
+          "title": "Uses",
+          "type": "string"
+        }
+      },
+      "required": [
+        "id",
+        "uses"
+      ],
+      "title": "OperatorTask",
+      "type": "object"
+    },
+    "TaskTemplate": {
+      "additionalProperties": true,
+      "description": "A reusable, partial task fragment merged into any task 
that ``extends`` it.\n\nThis holds arbitrary keys as-is. All validation is done 
on the merged task.",
+      "properties": {},
+      "title": "TaskTemplate",
+      "type": "object"
+    },
+    "TemplateRef": {
+      "additionalProperties": false,
+      "description": "A Jinja template.",
+      "properties": {
+        "$t": {
+          "title": "$T",
+          "type": "string"
+        }
+      },
+      "required": [
+        "$t"
+      ],
+      "title": "TemplateRef",
+      "type": "object"
+    },
+    "TimetableSchedule": {
+      "additionalProperties": false,
+      "description": "A constructed timetable.",
+      "properties": {
+        "uses": {
+          "title": "Uses",
+          "type": "string"
+        },
+        "with": {
+          "additionalProperties": {
+            "oneOf": [
+              {
+                "$ref": "#/$defs/XComRef"
+              },
+              {
+                "$ref": "#/$defs/TemplateRef"
+              },
+              {
+                "$ref": "#/$defs/ConstRef"
+              },
+              {
+                "$ref": "#/$defs/_Literal"
+              }
+            ]
+          },
+          "title": "With",
+          "type": "object"
+        }
+      },
+      "required": [
+        "uses"
+      ],
+      "title": "TimetableSchedule",
+      "type": "object"
+    },
+    "XComRef": {
+      "additionalProperties": false,
+      "description": "An upstream task's XCom output.",
+      "properties": {
+        "$x": {
+          "anyOf": [
+            {
+              "type": "string"
+            },
+            {
+              "$ref": "#/$defs/XComTarget"
+            }
+          ],
+          "title": "$X"
+        }
+      },
+      "required": [
+        "$x"
+      ],
+      "title": "XComRef",
+      "type": "object"
+    },
+    "XComTarget": {
+      "additionalProperties": false,
+      "description": "The object form of an XCom reference.",
+      "properties": {
+        "task": {
+          "title": "Task",
+          "type": "string"
+        },
+        "key": {
+          "anyOf": [
+            {
+              "type": "string"
+            },
+            {
+              "type": "null"
+            }
+          ],
+          "default": null,
+          "title": "Key"
+        }
+      },
+      "required": [
+        "task"
+      ],
+      "title": "XComTarget",
+      "type": "object"
+    },
+    "_Literal": {
+      "description": "A literal value.\n\nContainers recurse so nested markers 
are still resolved; scalars pass\nthrough.",
+      "title": "_Literal"
+    }
+  },
+  "additionalProperties": true,
+  "description": "A Dag represented by one YAML document.\n\nExtra top-level 
keys are Dag-level arguments.",
+  "properties": {
+    "$schema": {
+      "title": "$Schema",
+      "type": "string"
+    },
+    "dag_id": {
+      "title": "Dag Id",
+      "type": "string"
+    },
+    "schedule": {
+      "anyOf": [
+        {
+          "type": "string"
+        },
+        {
+          "$ref": "#/$defs/TimetableSchedule"
+        },
+        {
+          "type": "null"
+        }
+      ],
+      "default": null,
+      "title": "Schedule"
+    },
+    "templates": {
+      "additionalProperties": {
+        "$ref": "#/$defs/TaskTemplate"
+      },
+      "title": "Templates",
+      "type": "object"
+    },
+    "tasks": {
+      "items": {
+        "oneOf": [

Review Comment:
   This schema rejects the PR's own `example_dag.yaml`. Running 
`jsonschema.Draft202012Validator` over each document gives 2 errors per 
document. The `oneOf` branches overlap in two places. Here, `CodeTask` only 
requires `id` and allows extra keys, so every `uses:` task matches both 
branches. In the value unions (line 27 and its three copies), `_Literal` dumps 
as `{}`, so `{$x: up}`, `{$t: ...}` and `{$const: ...}` each match two branches 
as well. The reverse also happens: `{id: t}` with neither `uses` nor `run` 
passes, and `$xcom`/`$template` only get through via `_Literal`, because 
`by_alias=True` emits just the first alias.
   
   Could the branches be made disjoint? That would mean `CodeTask` requiring 
`run` and excluding `uses`, and `_Literal` excluding single-key marker objects, 
with both spellings listed. `test_model_json_schema_describes_markers` could 
then validate the fixture against schema.json; `jsonschema` is already a 
task-sdk dependency.



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/parser.py:
##########
@@ -0,0 +1,72 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+Parse a YAML/JSON DAG file into validated :class:`~.models.DagDocument` 
objects.
+
+The pydantic models (`.models`) are the schema; this module is the thin I/O + 
versioning
+layer around them: it lazily streams each document of a multi-document 
YAML/JSON stream (a string or a
+readable file object),
+resolves its `$schema` version against the Cadwyn bundle (migrating an 
older-pinned document
+up to head), validates it, and turns pydantic/YAML failures into a readable
+:class:`YamlDagParseError` at the offending document. Parsing is lazy, so 
errors surface as
+the iterator is consumed. It is pure (yaml + pydantic + cadwyn; no `airflow` 
import), so it is
+unit-testable without a runtime.
+"""
+
+from __future__ import annotations
+
+import collections.abc
+from typing import TYPE_CHECKING, Any
+
+import yaml
+from pydantic import ValidationError
+
+from airflow.sdk.importers.yaml_importer import migrator
+from airflow.sdk.importers.yaml_importer.models import DagDocument
+
+if TYPE_CHECKING:
+    from collections.abc import Iterator
+    from typing import IO
+
+
+class YamlDagParseError(ValueError):
+    """A document could not be parsed or did not conform to the format."""
+
+
+def _resolve_and_migrate(raw: dict[str, Any], *, source: str) -> dict[str, 
Any]:
+    """Resolve the ``$schema`` version and migrate to the head shape."""
+    if not isinstance(raw, collections.abc.Mapping):
+        raise YamlDagParseError(f"{source}: a DAG document must be a mapping, 
got {type(raw).__name__}")
+    if not raw.get("$schema"):
+        raise YamlDagParseError(f"{source}: missing required key '$schema'")
+    return migrator.get_migrator().resolve_and_migrate(raw, source=source)
+
+
+def parse_documents(stream: str | IO[str], *, source: str = "<string>") -> 
Iterator[DagDocument]:
+    """Parse each document of *stream* into a :class:`DagDocument`."""
+    try:
+        for pos, raw in enumerate(yaml.safe_load_all(stream)):

Review Comment:
   PyYAML keeps the last value for a repeated key, so a second `tasks:` block 
(or a second `with:` / `needs:` on one task) silently replaces the first. A 
document with `tasks: [...]` followed later by `tasks: []` parses as a Dag with 
no tasks. For a hand-edited format, could this use a `SafeLoader` subclass 
whose `construct_mapping` raises on a duplicate key? A `ConstructorError` is a 
`YAMLError`, so the handler below would already wrap it with the source.



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/models.py:
##########
@@ -0,0 +1,362 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+Pydantic models for the YAML/JSON DAG format.
+
+These are the single source of truth for the format; they validate an authored
+document (rich Pydantic errors) and back the published JSON Schema. Only the
+format skeleton and value grammar only is described; operator ``with:`` args,
+Dag- and task-level attributes pass through untyped because the Python
+constructors validate those at import time.
+"""
+
+from __future__ import annotations
+
+import itertools
+from typing import Annotated, Any
+
+from pydantic import (
+    AliasChoices,
+    BaseModel,
+    ConfigDict,
+    Discriminator,
+    Field,
+    RootModel,
+    Tag,
+    TypeAdapter,
+    model_validator,
+)
+
+XCOM_KEYS = ("$x", "$xcom")
+TEMPLATE_KEYS = ("$t", "$template")
+CONST_KEY = "$const"
+
+# Keys deferred to a later edition...
+DAG_CALLBACK_KEYS = {"on_success_callback", "on_failure_callback", 
"sla_miss_callback"}
+TASK_CALLBACK_KEYS = {
+    "on_success_callback",
+    "on_failure_callback",
+    "on_retry_callback",
+    "on_execute_callback",
+    "on_skipped_callback",
+    "sla_miss_callback",
+}
+TASK_ASSET_IO_KEYS = {"inlets", "outlets"}
+
+# Structural keys the format owns (everything else is pass-through to Python).
+_DAG_STRUCTURAL = {"$schema", "dag_id", "schedule", "templates", "tasks"}
+
+
+class XComTarget(BaseModel):
+    """The object form of an XCom reference."""
+
+    task: str
+    key: str | None = None
+
+    model_config = ConfigDict(extra="forbid")
+
+
+class XComRef(BaseModel):
+    """An upstream task's XCom output."""
+
+    target: str | XComTarget = Field(
+        serialization_alias="$x",
+        validation_alias=AliasChoices(*XCOM_KEYS),
+    )
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+class TemplateRef(BaseModel):
+    """A Jinja template."""
+
+    source: str = Field(serialization_alias="$t", 
validation_alias=AliasChoices(*TEMPLATE_KEYS))
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+class ConstRef(BaseModel):
+    """Force a literal; the value is taken verbatim."""
+
+    value: Any = Field(serialization_alias="$const", 
validation_alias=AliasChoices(CONST_KEY))
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+def _value_discriminator(v: Any) -> str:
+    """
+    Route a raw value to a marker branch, or to ``literal``.
+
+    A reserved ``$``-marker is recognised only as the sole key of an object, so
+    a multi-key dict that merely contains ``$x`` is a plain literal dict.
+    """
+    if isinstance(v, dict) and len(v) == 1:
+        (key,) = v
+        if key in XCOM_KEYS:
+            return "xcom"
+        if key in TEMPLATE_KEYS:
+            return "template"
+        if key == CONST_KEY:
+            return "const"
+    # Already-parsed marker instances (revalidation) route to themselves.
+    if isinstance(v, XComRef):
+        return "xcom"
+    if isinstance(v, TemplateRef):
+        return "template"
+    if isinstance(v, ConstRef):
+        return "const"
+    return "literal"
+
+
+class _Literal(RootModel[Any]):
+    """
+    A literal value.
+
+    Containers recurse so nested markers are still resolved; scalars pass
+    through.
+    """
+
+    root: Any

Review Comment:
   Because `_recurse` calls `_ValueAdapter.validate_python` inside a 
before-validator, errors in nested values lose their path. `run: {cfg: {u: {$x: 
5}}}` reports the location as `run.cfg.literal.xcom.$x`, with no `u`, so in a 
deep config the author can't tell which key is wrong. Typing the root as a 
recursive union (`dict[str, Value] | list[Value] | str | int | float | bool | 
None`) would let pydantic keep the full location.



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/models.py:
##########
@@ -0,0 +1,362 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+Pydantic models for the YAML/JSON DAG format.
+
+These are the single source of truth for the format; they validate an authored
+document (rich Pydantic errors) and back the published JSON Schema. Only the
+format skeleton and value grammar only is described; operator ``with:`` args,
+Dag- and task-level attributes pass through untyped because the Python
+constructors validate those at import time.
+"""
+
+from __future__ import annotations
+
+import itertools
+from typing import Annotated, Any
+
+from pydantic import (
+    AliasChoices,
+    BaseModel,
+    ConfigDict,
+    Discriminator,
+    Field,
+    RootModel,
+    Tag,
+    TypeAdapter,
+    model_validator,
+)
+
+XCOM_KEYS = ("$x", "$xcom")
+TEMPLATE_KEYS = ("$t", "$template")
+CONST_KEY = "$const"
+
+# Keys deferred to a later edition...
+DAG_CALLBACK_KEYS = {"on_success_callback", "on_failure_callback", 
"sla_miss_callback"}
+TASK_CALLBACK_KEYS = {
+    "on_success_callback",
+    "on_failure_callback",
+    "on_retry_callback",
+    "on_execute_callback",
+    "on_skipped_callback",
+    "sla_miss_callback",
+}
+TASK_ASSET_IO_KEYS = {"inlets", "outlets"}
+
+# Structural keys the format owns (everything else is pass-through to Python).
+_DAG_STRUCTURAL = {"$schema", "dag_id", "schedule", "templates", "tasks"}
+
+
+class XComTarget(BaseModel):
+    """The object form of an XCom reference."""
+
+    task: str
+    key: str | None = None
+
+    model_config = ConfigDict(extra="forbid")
+
+
+class XComRef(BaseModel):
+    """An upstream task's XCom output."""
+
+    target: str | XComTarget = Field(
+        serialization_alias="$x",
+        validation_alias=AliasChoices(*XCOM_KEYS),
+    )
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+class TemplateRef(BaseModel):
+    """A Jinja template."""
+
+    source: str = Field(serialization_alias="$t", 
validation_alias=AliasChoices(*TEMPLATE_KEYS))
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+class ConstRef(BaseModel):
+    """Force a literal; the value is taken verbatim."""
+
+    value: Any = Field(serialization_alias="$const", 
validation_alias=AliasChoices(CONST_KEY))
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+def _value_discriminator(v: Any) -> str:
+    """
+    Route a raw value to a marker branch, or to ``literal``.
+
+    A reserved ``$``-marker is recognised only as the sole key of an object, so
+    a multi-key dict that merely contains ``$x`` is a plain literal dict.
+    """
+    if isinstance(v, dict) and len(v) == 1:
+        (key,) = v
+        if key in XCOM_KEYS:
+            return "xcom"
+        if key in TEMPLATE_KEYS:
+            return "template"
+        if key == CONST_KEY:
+            return "const"
+    # Already-parsed marker instances (revalidation) route to themselves.
+    if isinstance(v, XComRef):
+        return "xcom"
+    if isinstance(v, TemplateRef):
+        return "template"
+    if isinstance(v, ConstRef):
+        return "const"
+    return "literal"
+
+
+class _Literal(RootModel[Any]):
+    """
+    A literal value.
+
+    Containers recurse so nested markers are still resolved; scalars pass
+    through.
+    """
+
+    root: Any
+
+    @model_validator(mode="before")
+    @classmethod
+    def _recurse(cls, v: Any) -> Any:
+        if isinstance(v, dict):
+            return {k: _ValueAdapter.validate_python(item) for k, item in 
v.items()}
+        if isinstance(v, list):
+            return [_ValueAdapter.validate_python(item) for item in v]
+        return v
+
+
+Value = Annotated[
+    Annotated[XComRef, Tag("xcom")]
+    | Annotated[TemplateRef, Tag("template")]
+    | Annotated[ConstRef, Tag("const")]
+    | Annotated[_Literal, Tag("literal")],
+    Discriminator(_value_discriminator),
+]
+
+_ValueAdapter: TypeAdapter[Any] = TypeAdapter(Value)
+
+
+class _TaskBase(BaseModel):
+    """
+    Fields common to both ``use:`` and ``run:`` task kinds.
+
+    Unknown keys are BaseOperator arguments (pass-through).
+    """
+
+    id_: str = Field(alias="id")
+    needs: list[str] = Field(default_factory=list)
+    extends: list[str] = Field(default_factory=list)
+    with_: dict[str, Value] = Field(default_factory=dict, alias="with")
+
+    model_config = ConfigDict(extra="allow", populate_by_name=True)
+
+    @model_validator(mode="before")
+    @classmethod
+    def _check_task(cls, data: Any) -> Any:
+        if isinstance(data, dict):
+            for k in TASK_ASSET_IO_KEYS:
+                if k in data:
+                    raise ValueError(f"{k!r} not yet implemented")
+            for k in TASK_CALLBACK_KEYS:
+                if k in data:
+                    raise ValueError(f"{k!r} (callback) not yet implemented")
+            bodies = [k for k in ("uses", "run") if k in data]
+            if len(bodies) != 1:
+                raise ValueError(
+                    f"task {data.get('id')!r} must have exactly one of 'uses' 
or 'run'"
+                    + (f"; got {bodies}" if bodies else "")
+                )
+        return data
+
+
+class OperatorTask(_TaskBase):
+    """
+    A ready-made operator.
+
+    This should have ``uses`` (import path) + ``with``.
+    """
+
+    uses: str
+
+
+class CodeTask(_TaskBase):
+    """
+    A task running custom code.
+
+    This should have ``run`` holding function arguments, and optionally
+    ``queue`` to route it.
+    """
+
+    run: dict[str, Value] = Field(default_factory=dict)
+
+
+def _task_discriminator(v: Any) -> str:
+    if isinstance(v, OperatorTask):
+        return "operator"
+    if isinstance(v, CodeTask):
+        return "code"
+    if isinstance(v, dict) and "uses" in v:
+        return "operator"
+    return "code"
+
+
+Task = Annotated[
+    Annotated[OperatorTask, Tag("operator")] | Annotated[CodeTask, 
Tag("code")],
+    Discriminator(_task_discriminator),
+]
+
+
+class TaskTemplate(BaseModel):
+    """
+    A reusable, partial task fragment merged into any task that ``extends`` it.
+
+    This holds arbitrary keys as-is. All validation is done on the merged task.
+    """
+
+    model_config = ConfigDict(extra="allow")
+
+
+class TimetableSchedule(BaseModel):
+    """A constructed timetable."""
+
+    uses: str
+    with_: dict[str, Value] = Field(default_factory=dict, alias="with")
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+Schedule = str | None | TimetableSchedule
+
+
+def _merge_task(base: dict, over: dict) -> dict:
+    """
+    Shallow merge *over* onto *base*.
+
+    Note that ``with`` and ``run`` need to merge one level deep. All other
+    fields simply replace.
+    """
+    out = base.copy()
+    for k, v in over.items():
+        if k in ("with", "run") and isinstance(v, dict) and 
isinstance(out.get(k), dict):
+            merged = dict(out[k])
+            merged.update(v)
+            out[k] = merged
+        else:
+            out[k] = v
+    return out
+
+
+def _expand_template(name: str, templates: dict[str, Any], _seen: tuple = ()) 
-> dict:
+    if name in _seen:
+        raise ValueError(f"template cycle via {name!r}")
+    if name not in templates:
+        raise ValueError(f"unknown template {name!r}")
+    spec = templates[name]
+    if not isinstance(spec, dict):
+        raise ValueError(f"template {name!r} must be a mapping")
+    acc: dict = {}
+    for parent in spec.get("extends", []) or []:
+        acc = _merge_task(acc, _expand_template(parent, templates, _seen + 
(name,)))
+    own = {k: v for k, v in spec.items() if k != "extends"}
+    return _merge_task(acc, own)
+
+
+def _iter_xcom_task_ids(value: Any):
+    """Yield task ids referenced by any XComRef inside a resolved value."""
+    if isinstance(value, XComRef):
+        yield value.target if isinstance(value.target, str) else 
value.target.task
+    elif isinstance(value, _Literal):
+        yield from _iter_xcom_task_ids(value.root)
+    elif isinstance(value, dict):
+        for v in value.values():
+            yield from _iter_xcom_task_ids(v)
+    elif isinstance(value, list):
+        for v in value:
+            yield from _iter_xcom_task_ids(v)
+
+
+class DagDocument(BaseModel):
+    """
+    A Dag represented by one YAML document.
+
+    Extra top-level keys are Dag-level arguments.
+    """
+
+    schema_: str = Field(alias="$schema")
+    dag_id: str
+    schedule: Schedule = None
+    templates: dict[str, TaskTemplate] = Field(default_factory=dict)
+    tasks: list[Task] = Field(default_factory=list)
+
+    model_config = ConfigDict(extra="allow", populate_by_name=True)
+
+    @model_validator(mode="before")
+    @classmethod
+    def _reject_deferred(cls, data: Any) -> Any:
+        if isinstance(data, dict):
+            if "default_args" in data:
+                raise ValueError("default_args is not supported; use 
'templates' with 'extends' instead")
+            for k in DAG_CALLBACK_KEYS:
+                if k in data:
+                    raise ValueError(f"Dag-level {k!r} is not implemented yet")
+        return data
+
+    @model_validator(mode="before")
+    @classmethod
+    def _resolve_extends(cls, data: Any) -> Any:
+        """Merge in each task's ``extends`` templates before task 
validation."""
+        if not isinstance(data, dict):
+            return data
+        templates = data.get("templates")
+        if templates is None:
+            templates = {}
+        tasks = data.get("tasks")
+        # Only merge when the shapes are as expected. A malformed input is left
+        # untouched, so normal field validation reports it as a 
ValidationError.
+        if not isinstance(templates, dict) or not isinstance(tasks, list):
+            return data
+        resolved = []
+        for task in tasks:
+            extends = task.get("extends") if isinstance(task, dict) else None
+            if isinstance(extends, list) and extends:
+                merged: dict = {}
+                for name in extends:
+                    merged = _merge_task(merged, _expand_template(name, 
templates))
+                resolved.append(_merge_task(merged, {k: v for k, v in 
task.items() if k != "extends"}))
+            else:
+                resolved.append(task)
+        return {**data, "tasks": resolved}
+
+    @model_validator(mode="after")
+    def _check_refs(self):
+        """Check ``needs`` and XCom targets are present in this Dag."""
+        ids = {t.id_ for t in self.tasks}
+        for t in self.tasks:
+            refs = set(t.needs)
+            for v in itertools.chain(t.with_.values(), getattr(t, "run", 
{}).values()):
+                refs.update(_iter_xcom_task_ids(v))
+            if missing := refs - ids:
+                raise ValueError(f"task {t.id_!r} references unknown task(s): 
{sorted(missing)}")
+        return self
+
+    @property
+    def dag_attributes(self) -> dict[str, Any]:
+        """Dag attributes beyond structural keys."""
+        return {k: v for k, v in (self.__pydantic_extra__ or {}).items() if k 
not in _DAG_STRUCTURAL}

Review Comment:
   Declared fields never end up in `__pydantic_extra__`, so this filter never 
removes anything: `$schema`, `dag_id` and the rest are already on their fields. 
`return dict(self.__pydantic_extra__ or {})` would do, and `_DAG_STRUCTURAL` 
could go.



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/migrator.py:
##########
@@ -0,0 +1,132 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+`$schema` version resolution and forward migration to the head shape.
+
+A document pinned to an older `$schema` version is migrated up to the current 
head shape
+*before* validation, by walking the Cadwyn version bundle and applying each 
`VersionChange`'s
+forward request converters 
(`@convert_request_to_next_version_for(DagDocument)`), oldest to
+newest. This mirrors the exec-API / supervisor-schema `SchemaVersionMigrator`, 
trimmed to the
+one direction a document format needs (author's version -> head), with the 
bundle and its
+version list encapsulated in the migrator rather than exposed as module-level 
helpers.
+
+With a single published version there is nothing to migrate; the machinery is 
here so a future
+version bump only has to add a `VersionChange` carrying a converter — no 
parser changes.
+"""
+
+from __future__ import annotations
+
+import copy
+import functools
+import re
+import warnings
+from typing import TYPE_CHECKING, Any
+
+import attrs
+
+from airflow.sdk.importers.yaml_importer.models import DagDocument
+from airflow.sdk.importers.yaml_importer.versions import get_bundle
+
+if TYPE_CHECKING:
+    from cadwyn import VersionBundle
+
+_SCHEMA_URL = "https://airflow.apache.org/schemas/dag/{date}.json";
+_VERSION_RE = re.compile(r"\d{4}-\d{2}-\d{2}")
+
+
+def schema_url(date: str) -> str:
+    """Build the canonical ``$schema`` URL for a published version *date*."""
+    return _SCHEMA_URL.format(date=date)
+
+
+def version_from_schema(schema: str) -> str | None:
+    """Read the version (``YYYY-MM-DD``) token out of a ``$schema`` URL, 
ignoring the host."""
+    match = _VERSION_RE.search(schema or "")
+    return match.group(0) if match else None
+
+
[email protected]
+class _RequestInfo:
+    """
+    Duck-type stand-in for Cadwyn's ``RequestInfo``.
+
+    ``cadwyn.structure.data._AlterDataInstruction.__call__`` only reads
+    and writes ``info.body``; the by-schema transformers we drive never
+    touch FastAPI's Request/Response. Passing this minimal object lets
+    us run cadwyn's migrations from a pure in-process code path with no
+    HTTP stack.
+    """
+
+    body: dict[str, Any]
+
+
[email protected]
+class DagDocumentMigrator:
+    """YAML Dag document migrator; pins each document to its ``$schema`` 
version."""
+
+    _bundle: VersionBundle
+
+    def resolve_and_migrate(self, body: dict[str, Any], *, source: str) -> 
dict[str, Any]:
+        """
+        Resolve *body*'s ``$schema`` version and migrate it to the head shape.
+
+        The version token is read from the ``$schema`` URL (the host is 
ignored). It is an
+        exact pin: a known version uses its own ruleset; an unknown one (e.g. 
newer than this
+        importer) resolves to the latest ruleset with a warning, never a hard 
failure.
+
+        :return: The migrated result. If *body* is already at head, it is
+            returned as-is; otherwise a migrated copy is returned.
+        """
+        known = [v.value for v in self._bundle.versions]  # newest-first
+        date = version_from_schema(body["$schema"])
+        if date in known:
+            source_version = date
+        else:
+            source_version = known[0]  # unknown version -> latest ruleset we 
have
+            warnings.warn(
+                f"{source}: $schema version {date!r} is not a known version 
{known}; "
+                f"using the latest ruleset {source_version!r}",
+                stacklevel=2,
+            )
+        return self._migrate_to_head(body, source_version)
+
+    def _migrate_to_head(self, body: dict[str, Any], source_version: str) -> 
dict[str, Any]:
+        """
+        Migrate a raw document *body* from *source_version* to the head shape.
+
+        This applies the forward request converters for :class:`DagDocument`
+        from the version after *source_version* through head, in order.
+
+        :return: The migrated result. If *body* is already at head, it is
+            returned as-is; otherwise a migrated copy is returned.
+        """
+        if source_version == self._bundle.versions[0].value:
+            return body
+        info = _RequestInfo(copy.deepcopy(dict(body)))
+        for version in self._bundle.reversed_versions:
+            if version.value <= source_version:
+                continue
+            for change in version.changes:
+                for instruction in 
change.alter_request_by_schema_instructions.get(DagDocument, ()):

Review Comment:
   Only `alter_request_by_schema_instructions[DagDocument]` is read here, so 
every other instruction on a `VersionChange` is dropped without an error. 
`schema(DagDocument).field("schedule").had(name="schedule_interval")`, the form 
`execution_time/schema/AGENTS.md` teaches, leaves a 2026-10-30 document's 
`schedule_interval` unmigrated. It then passes through as a Dag kwarg with 
`schedule=None`. Converters keyed on `OperatorTask` or `CodeTask` are skipped 
the same way. Could the migrator reject a bundle that carries anything else, or 
at least have the docstring say only whole-document `DagDocument` request 
converters are applied?



##########
task-sdk/tests/task_sdk/importers/yaml_importer/test_parser.py:
##########
@@ -0,0 +1,328 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""Unit tests for the YAML DAG format parser 
(airflow.sdk.importers.yaml_importer)."""
+
+from __future__ import annotations
+
+import io
+import textwrap
+
+import pytest
+
+from airflow.sdk.importers.yaml_importer import YamlDagParseError, 
parse_documents
+from airflow.sdk.importers.yaml_importer.models import (
+    CodeTask,
+    ConstRef,
+    DagDocument,
+    OperatorTask,
+    TemplateRef,
+    TimetableSchedule,
+    XComRef,
+    XComTarget,
+    _Literal,
+)
+
+HEAD = "2026-10-30"
+HEAD_SCHEMA = f"https://airflow.apache.org/schemas/dag/{HEAD}.json";
+
+
+def _one(body: str):
+    return next(iter(parse_documents(f"$schema: {HEAD_SCHEMA}\ndag_id: 
d\n{textwrap.dedent(body)}")))
+
+
+def _task(task_yaml: str):
+    return _one("tasks:\n" + textwrap.indent(textwrap.dedent(task_yaml), "  
")).tasks[0]
+
+
+# --------------------------------------------------------------------------- 
value grammar
+def test_literal_by_default():
+    t = _task("- {id: t, run: {region: 'eu', n: 5, flag: true}}")
+    assert isinstance(t.run["region"], _Literal)
+    assert t.run["region"].root == "eu"
+    assert t.run["n"].root == 5
+    assert t.run["flag"].root is True
+
+
+def test_xcom_short_long_and_keyed():
+    doc = _one(
+        "tasks:\n  - {id: up, run: {}}\n  - {id: t, run: {a: {$x: up}, b: 
{$xcom: up}, c: {$x: {task: up, key: rows}}}}"
+    )
+    t = {x.id_: x for x in doc.tasks}["t"]
+    assert isinstance(t.run["a"], XComRef)
+    assert t.run["a"].target == "up"
+    assert isinstance(t.run["b"], XComRef)
+    assert t.run["b"].target == "up"
+    assert isinstance(t.run["c"].target, XComTarget)
+    assert t.run["c"].target.key == "rows"
+
+
+def test_template_marker():
+    t = _task("- {id: t, uses: X, with: {k: {$t: 'v-{{ ds }}'}, k2: 
{$template: 'x'}}}")
+    assert isinstance(t.with_["k"], TemplateRef)
+    assert "{{ ds }}" in t.with_["k"].source
+    assert isinstance(t.with_["k2"], TemplateRef)
+
+
+def test_const_is_verbatim_not_recursed():
+    t = _task("- {id: t, run: {a: {$const: {$x: not-a-ref}}}}")
+    assert isinstance(t.run["a"], ConstRef)
+    assert t.run["a"].value == {"$x": "not-a-ref"}  # inner marker NOT 
interpreted
+
+
+def test_marker_only_as_sole_key():
+    t = _task("- {id: t, run: {a: {$x: up, extra: 1}}}")  # two keys -> 
literal dict
+    assert isinstance(t.run["a"], _Literal)
+    assert set(t.run["a"].root) == {"$x", "extra"}
+
+
+def test_nested_marker_inside_literal_resolved():
+    t = _task("- {id: t, uses: X, with: {cfg: {url: {$t: '{{ ds }}'}, n: 1}}}")
+    inner = t.with_["cfg"].root
+    assert isinstance(inner["url"], TemplateRef)
+    assert inner["n"].root == 1
+
+
+# --------------------------------------------------------------------------- 
task bodies
+def test_operator_vs_code_discrimination():
+    assert isinstance(_task("- {id: t, uses: a.b.C, with: {x: 1}}"), 
OperatorTask)
+    assert isinstance(_task("- {id: t, run: {}}"), CodeTask)
+
+
+def test_no_body_is_error():
+    with pytest.raises(YamlDagParseError, match="exactly one"):
+        _task("- {id: t}")
+
+
+def test_both_bodies_is_error():
+    with pytest.raises(YamlDagParseError, match="exactly one"):
+        _task("- {id: t, uses: X, run: {}}")
+
+
+# --------------------------------------------------------------------------- 
deferred keys
+def test_default_args_rejected():
+    with pytest.raises(YamlDagParseError, match="default_args"):
+        _one("default_args: {retries: 1}\ntasks: []")
+
+
+def test_callbacks_rejected():
+    with pytest.raises(YamlDagParseError, match="callback"):
+        _one("on_failure_callback: cb\ntasks: []")
+    with pytest.raises(YamlDagParseError, match="callback"):
+        _task("- {id: t, run: {}, on_success_callback: cb}")
+
+
+def test_inlets_outlets_rejected():
+    with pytest.raises(YamlDagParseError, match="not yet implemented"):
+        _task("- {id: t, run: {}, outlets: [a]}")
+
+
+# --------------------------------------------------------------------------- 
templates / extends
+def test_extends_merges_shallow_task_wins():
+    doc = _one(
+        """
+        templates:
+          retryable: {retries: 2, retry_delay: '5m'}
+          load:
+            extends: [retryable]
+            uses: a.b.S3ToRedshiftOperator
+            with: {schema: public, conn_id: warehouse}
+        tasks:
+          - id: t
+            extends: [load]
+            with: {table: sales, conn_id: override}
+        """
+    )
+    t = doc.tasks[0]
+    assert isinstance(t, OperatorTask)
+    assert t.uses.endswith("S3ToRedshiftOperator")
+    assert t.__pydantic_extra__["retries"] == 2  # from retryable
+    assert isinstance(t.with_["schema"], _Literal)  # inherited
+    assert t.with_["table"].root == "sales"  # own
+    assert t.with_["conn_id"].root == "override"  # one-level with-merge, task 
wins
+
+
+def test_unknown_template_errors():
+    with pytest.raises(YamlDagParseError, match="unknown template"):
+        _one("tasks:\n  - {id: t, extends: [ghost], run: {}}")
+
+
+def test_template_cycle_errors():
+    with pytest.raises(YamlDagParseError, match="cycle"):
+        _one("templates: {a: {extends: [b]}, b: {extends: [a]}}\ntasks: [{id: 
t, extends: [a], run: {}}]")
+
+
+# --------------------------------------------------------------------------- 
references / edges
+def test_unknown_needs_rejected():
+    with pytest.raises(YamlDagParseError, match="unknown task"):
+        _one("tasks:\n  - {id: t, run: {}, needs: [ghost]}")
+
+
+def test_unknown_xcom_target_rejected():

Review Comment:
   Only a top-level `$x` is covered. Deleting the list branch of 
`_Literal._recurse`, or the `_Literal` or list branch of `_iter_xcom_task_ids`, 
keeps the suite green, so nested references like `run: {a: [{$x: ghost}]}`, 
`run: {a: {b: {$x: ghost}}}` and `with: {a: {$x: ghost}}` are unpinned. 
Parametrizing this test over those shapes would catch it. The parser's own 
error paths (a non-mapping document, invalid YAML, the `[doc N]` label) have no 
tests either.



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/models.py:
##########
@@ -0,0 +1,362 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+Pydantic models for the YAML/JSON DAG format.
+
+These are the single source of truth for the format; they validate an authored
+document (rich Pydantic errors) and back the published JSON Schema. Only the
+format skeleton and value grammar only is described; operator ``with:`` args,
+Dag- and task-level attributes pass through untyped because the Python
+constructors validate those at import time.
+"""
+
+from __future__ import annotations
+
+import itertools
+from typing import Annotated, Any
+
+from pydantic import (
+    AliasChoices,
+    BaseModel,
+    ConfigDict,
+    Discriminator,
+    Field,
+    RootModel,
+    Tag,
+    TypeAdapter,
+    model_validator,
+)
+
+XCOM_KEYS = ("$x", "$xcom")
+TEMPLATE_KEYS = ("$t", "$template")
+CONST_KEY = "$const"
+
+# Keys deferred to a later edition...
+DAG_CALLBACK_KEYS = {"on_success_callback", "on_failure_callback", 
"sla_miss_callback"}
+TASK_CALLBACK_KEYS = {
+    "on_success_callback",
+    "on_failure_callback",
+    "on_retry_callback",
+    "on_execute_callback",
+    "on_skipped_callback",
+    "sla_miss_callback",
+}
+TASK_ASSET_IO_KEYS = {"inlets", "outlets"}
+
+# Structural keys the format owns (everything else is pass-through to Python).
+_DAG_STRUCTURAL = {"$schema", "dag_id", "schedule", "templates", "tasks"}
+
+
+class XComTarget(BaseModel):
+    """The object form of an XCom reference."""
+
+    task: str
+    key: str | None = None
+
+    model_config = ConfigDict(extra="forbid")
+
+
+class XComRef(BaseModel):
+    """An upstream task's XCom output."""
+
+    target: str | XComTarget = Field(
+        serialization_alias="$x",
+        validation_alias=AliasChoices(*XCOM_KEYS),
+    )
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+class TemplateRef(BaseModel):
+    """A Jinja template."""
+
+    source: str = Field(serialization_alias="$t", 
validation_alias=AliasChoices(*TEMPLATE_KEYS))
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+class ConstRef(BaseModel):
+    """Force a literal; the value is taken verbatim."""
+
+    value: Any = Field(serialization_alias="$const", 
validation_alias=AliasChoices(CONST_KEY))
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+def _value_discriminator(v: Any) -> str:
+    """
+    Route a raw value to a marker branch, or to ``literal``.
+
+    A reserved ``$``-marker is recognised only as the sole key of an object, so
+    a multi-key dict that merely contains ``$x`` is a plain literal dict.
+    """
+    if isinstance(v, dict) and len(v) == 1:
+        (key,) = v
+        if key in XCOM_KEYS:
+            return "xcom"
+        if key in TEMPLATE_KEYS:
+            return "template"
+        if key == CONST_KEY:
+            return "const"
+    # Already-parsed marker instances (revalidation) route to themselves.
+    if isinstance(v, XComRef):
+        return "xcom"
+    if isinstance(v, TemplateRef):
+        return "template"
+    if isinstance(v, ConstRef):
+        return "const"
+    return "literal"
+
+
+class _Literal(RootModel[Any]):
+    """
+    A literal value.
+
+    Containers recurse so nested markers are still resolved; scalars pass
+    through.
+    """
+
+    root: Any
+
+    @model_validator(mode="before")
+    @classmethod
+    def _recurse(cls, v: Any) -> Any:
+        if isinstance(v, dict):
+            return {k: _ValueAdapter.validate_python(item) for k, item in 
v.items()}
+        if isinstance(v, list):
+            return [_ValueAdapter.validate_python(item) for item in v]
+        return v
+
+
+Value = Annotated[
+    Annotated[XComRef, Tag("xcom")]
+    | Annotated[TemplateRef, Tag("template")]
+    | Annotated[ConstRef, Tag("const")]
+    | Annotated[_Literal, Tag("literal")],
+    Discriminator(_value_discriminator),
+]
+
+_ValueAdapter: TypeAdapter[Any] = TypeAdapter(Value)
+
+
+class _TaskBase(BaseModel):
+    """
+    Fields common to both ``use:`` and ``run:`` task kinds.
+
+    Unknown keys are BaseOperator arguments (pass-through).
+    """
+
+    id_: str = Field(alias="id")
+    needs: list[str] = Field(default_factory=list)
+    extends: list[str] = Field(default_factory=list)
+    with_: dict[str, Value] = Field(default_factory=dict, alias="with")
+
+    model_config = ConfigDict(extra="allow", populate_by_name=True)
+
+    @model_validator(mode="before")
+    @classmethod
+    def _check_task(cls, data: Any) -> Any:
+        if isinstance(data, dict):
+            for k in TASK_ASSET_IO_KEYS:
+                if k in data:
+                    raise ValueError(f"{k!r} not yet implemented")
+            for k in TASK_CALLBACK_KEYS:
+                if k in data:
+                    raise ValueError(f"{k!r} (callback) not yet implemented")
+            bodies = [k for k in ("uses", "run") if k in data]
+            if len(bodies) != 1:
+                raise ValueError(
+                    f"task {data.get('id')!r} must have exactly one of 'uses' 
or 'run'"
+                    + (f"; got {bodies}" if bodies else "")
+                )
+        return data
+
+
+class OperatorTask(_TaskBase):
+    """
+    A ready-made operator.
+
+    This should have ``uses`` (import path) + ``with``.
+    """
+
+    uses: str
+
+
+class CodeTask(_TaskBase):
+    """
+    A task running custom code.
+
+    This should have ``run`` holding function arguments, and optionally
+    ``queue`` to route it.
+    """
+
+    run: dict[str, Value] = Field(default_factory=dict)
+
+
+def _task_discriminator(v: Any) -> str:
+    if isinstance(v, OperatorTask):
+        return "operator"
+    if isinstance(v, CodeTask):
+        return "code"
+    if isinstance(v, dict) and "uses" in v:
+        return "operator"
+    return "code"
+
+
+Task = Annotated[
+    Annotated[OperatorTask, Tag("operator")] | Annotated[CodeTask, 
Tag("code")],
+    Discriminator(_task_discriminator),
+]
+
+
+class TaskTemplate(BaseModel):
+    """
+    A reusable, partial task fragment merged into any task that ``extends`` it.
+
+    This holds arbitrary keys as-is. All validation is done on the merged task.
+    """
+
+    model_config = ConfigDict(extra="allow")
+
+
+class TimetableSchedule(BaseModel):
+    """A constructed timetable."""
+
+    uses: str
+    with_: dict[str, Value] = Field(default_factory=dict, alias="with")
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+Schedule = str | None | TimetableSchedule
+
+
+def _merge_task(base: dict, over: dict) -> dict:
+    """
+    Shallow merge *over* onto *base*.
+
+    Note that ``with`` and ``run`` need to merge one level deep. All other
+    fields simply replace.
+    """
+    out = base.copy()
+    for k, v in over.items():
+        if k in ("with", "run") and isinstance(v, dict) and 
isinstance(out.get(k), dict):
+            merged = dict(out[k])
+            merged.update(v)
+            out[k] = merged
+        else:
+            out[k] = v
+    return out
+
+
+def _expand_template(name: str, templates: dict[str, Any], _seen: tuple = ()) 
-> dict:
+    if name in _seen:
+        raise ValueError(f"template cycle via {name!r}")
+    if name not in templates:

Review Comment:
   A non-string entry such as `extends: [[base]]` or `extends: [{base: 1}]` 
raises `TypeError: unhashable type` here. Pydantic only converts 
`ValueError`/`AssertionError` from a validator, so it escapes `parse_documents` 
unwrapped. Separately, a template's own `extends: base` (a string, not a list) 
is iterated character by character on line 275. It fails with `unknown template 
'b'`, or silently works when the template name is one character long. Checking 
that `extends` is a list of strings and raising `ValueError` would cover both.



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/parser.py:
##########
@@ -0,0 +1,72 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+Parse a YAML/JSON DAG file into validated :class:`~.models.DagDocument` 
objects.
+
+The pydantic models (`.models`) are the schema; this module is the thin I/O + 
versioning
+layer around them: it lazily streams each document of a multi-document 
YAML/JSON stream (a string or a
+readable file object),
+resolves its `$schema` version against the Cadwyn bundle (migrating an 
older-pinned document
+up to head), validates it, and turns pydantic/YAML failures into a readable
+:class:`YamlDagParseError` at the offending document. Parsing is lazy, so 
errors surface as
+the iterator is consumed. It is pure (yaml + pydantic + cadwyn; no `airflow` 
import), so it is
+unit-testable without a runtime.
+"""
+
+from __future__ import annotations
+
+import collections.abc
+from typing import TYPE_CHECKING, Any
+
+import yaml
+from pydantic import ValidationError
+
+from airflow.sdk.importers.yaml_importer import migrator
+from airflow.sdk.importers.yaml_importer.models import DagDocument
+
+if TYPE_CHECKING:
+    from collections.abc import Iterator
+    from typing import IO
+
+
+class YamlDagParseError(ValueError):
+    """A document could not be parsed or did not conform to the format."""
+
+
+def _resolve_and_migrate(raw: dict[str, Any], *, source: str) -> dict[str, 
Any]:
+    """Resolve the ``$schema`` version and migrate to the head shape."""
+    if not isinstance(raw, collections.abc.Mapping):
+        raise YamlDagParseError(f"{source}: a DAG document must be a mapping, 
got {type(raw).__name__}")
+    if not raw.get("$schema"):

Review Comment:
   This only checks truthiness. `$schema: 2026-10-30` without quotes loads as a 
`datetime.date`, reaches `re.search` in `version_from_schema`, and escapes as 
`TypeError: expected string or bytes-like object, got 'datetime.date'` rather 
than `YamlDagParseError`. An int or a list does the same. Since this key 
replaced the date-valued `compatibility_date`, the unquoted date seems a likely 
slip. An `isinstance(..., str)` check here would keep the error contract.



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/models.py:
##########
@@ -0,0 +1,362 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+Pydantic models for the YAML/JSON DAG format.
+
+These are the single source of truth for the format; they validate an authored
+document (rich Pydantic errors) and back the published JSON Schema. Only the
+format skeleton and value grammar only is described; operator ``with:`` args,
+Dag- and task-level attributes pass through untyped because the Python
+constructors validate those at import time.
+"""
+
+from __future__ import annotations
+
+import itertools
+from typing import Annotated, Any
+
+from pydantic import (
+    AliasChoices,
+    BaseModel,
+    ConfigDict,
+    Discriminator,
+    Field,
+    RootModel,
+    Tag,
+    TypeAdapter,
+    model_validator,
+)
+
+XCOM_KEYS = ("$x", "$xcom")
+TEMPLATE_KEYS = ("$t", "$template")
+CONST_KEY = "$const"
+
+# Keys deferred to a later edition...
+DAG_CALLBACK_KEYS = {"on_success_callback", "on_failure_callback", 
"sla_miss_callback"}
+TASK_CALLBACK_KEYS = {
+    "on_success_callback",
+    "on_failure_callback",
+    "on_retry_callback",
+    "on_execute_callback",
+    "on_skipped_callback",
+    "sla_miss_callback",
+}
+TASK_ASSET_IO_KEYS = {"inlets", "outlets"}
+
+# Structural keys the format owns (everything else is pass-through to Python).
+_DAG_STRUCTURAL = {"$schema", "dag_id", "schedule", "templates", "tasks"}
+
+
+class XComTarget(BaseModel):
+    """The object form of an XCom reference."""
+
+    task: str
+    key: str | None = None
+
+    model_config = ConfigDict(extra="forbid")
+
+
+class XComRef(BaseModel):
+    """An upstream task's XCom output."""
+
+    target: str | XComTarget = Field(
+        serialization_alias="$x",
+        validation_alias=AliasChoices(*XCOM_KEYS),
+    )
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+class TemplateRef(BaseModel):
+    """A Jinja template."""
+
+    source: str = Field(serialization_alias="$t", 
validation_alias=AliasChoices(*TEMPLATE_KEYS))
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+class ConstRef(BaseModel):
+    """Force a literal; the value is taken verbatim."""
+
+    value: Any = Field(serialization_alias="$const", 
validation_alias=AliasChoices(CONST_KEY))
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+def _value_discriminator(v: Any) -> str:
+    """
+    Route a raw value to a marker branch, or to ``literal``.
+
+    A reserved ``$``-marker is recognised only as the sole key of an object, so
+    a multi-key dict that merely contains ``$x`` is a plain literal dict.
+    """
+    if isinstance(v, dict) and len(v) == 1:
+        (key,) = v
+        if key in XCOM_KEYS:
+            return "xcom"
+        if key in TEMPLATE_KEYS:
+            return "template"
+        if key == CONST_KEY:
+            return "const"
+    # Already-parsed marker instances (revalidation) route to themselves.
+    if isinstance(v, XComRef):
+        return "xcom"
+    if isinstance(v, TemplateRef):
+        return "template"
+    if isinstance(v, ConstRef):
+        return "const"
+    return "literal"
+
+
+class _Literal(RootModel[Any]):
+    """
+    A literal value.
+
+    Containers recurse so nested markers are still resolved; scalars pass
+    through.
+    """
+
+    root: Any
+
+    @model_validator(mode="before")
+    @classmethod
+    def _recurse(cls, v: Any) -> Any:
+        if isinstance(v, dict):
+            return {k: _ValueAdapter.validate_python(item) for k, item in 
v.items()}
+        if isinstance(v, list):
+            return [_ValueAdapter.validate_python(item) for item in v]
+        return v
+
+
+Value = Annotated[
+    Annotated[XComRef, Tag("xcom")]
+    | Annotated[TemplateRef, Tag("template")]
+    | Annotated[ConstRef, Tag("const")]
+    | Annotated[_Literal, Tag("literal")],
+    Discriminator(_value_discriminator),
+]
+
+_ValueAdapter: TypeAdapter[Any] = TypeAdapter(Value)
+
+
+class _TaskBase(BaseModel):
+    """
+    Fields common to both ``use:`` and ``run:`` task kinds.
+
+    Unknown keys are BaseOperator arguments (pass-through).
+    """
+
+    id_: str = Field(alias="id")
+    needs: list[str] = Field(default_factory=list)
+    extends: list[str] = Field(default_factory=list)
+    with_: dict[str, Value] = Field(default_factory=dict, alias="with")
+
+    model_config = ConfigDict(extra="allow", populate_by_name=True)

Review Comment:
   `populate_by_name=True` makes the Python field names part of the format. 
`{id_: a, run: {}}` parses as task `a`, and a `with_` key next to `with` 
silently lands in `__pydantic_extra__`, so it would be passed to the operator. 
Nothing in the PR constructs these models by field name. Could it be dropped 
here and on lines 81, 89, 97, 242 and 308? Otherwise the parser accepts 
documents that schema.json rejects.



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/models.py:
##########
@@ -0,0 +1,362 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+Pydantic models for the YAML/JSON DAG format.
+
+These are the single source of truth for the format; they validate an authored
+document (rich Pydantic errors) and back the published JSON Schema. Only the
+format skeleton and value grammar only is described; operator ``with:`` args,
+Dag- and task-level attributes pass through untyped because the Python
+constructors validate those at import time.
+"""
+
+from __future__ import annotations
+
+import itertools
+from typing import Annotated, Any
+
+from pydantic import (
+    AliasChoices,
+    BaseModel,
+    ConfigDict,
+    Discriminator,
+    Field,
+    RootModel,
+    Tag,
+    TypeAdapter,
+    model_validator,
+)
+
+XCOM_KEYS = ("$x", "$xcom")
+TEMPLATE_KEYS = ("$t", "$template")
+CONST_KEY = "$const"
+
+# Keys deferred to a later edition...
+DAG_CALLBACK_KEYS = {"on_success_callback", "on_failure_callback", 
"sla_miss_callback"}
+TASK_CALLBACK_KEYS = {
+    "on_success_callback",
+    "on_failure_callback",
+    "on_retry_callback",
+    "on_execute_callback",
+    "on_skipped_callback",
+    "sla_miss_callback",
+}
+TASK_ASSET_IO_KEYS = {"inlets", "outlets"}
+
+# Structural keys the format owns (everything else is pass-through to Python).
+_DAG_STRUCTURAL = {"$schema", "dag_id", "schedule", "templates", "tasks"}
+
+
+class XComTarget(BaseModel):
+    """The object form of an XCom reference."""
+
+    task: str
+    key: str | None = None
+
+    model_config = ConfigDict(extra="forbid")
+
+
+class XComRef(BaseModel):
+    """An upstream task's XCom output."""
+
+    target: str | XComTarget = Field(
+        serialization_alias="$x",
+        validation_alias=AliasChoices(*XCOM_KEYS),
+    )
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+class TemplateRef(BaseModel):
+    """A Jinja template."""
+
+    source: str = Field(serialization_alias="$t", 
validation_alias=AliasChoices(*TEMPLATE_KEYS))
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+class ConstRef(BaseModel):
+    """Force a literal; the value is taken verbatim."""
+
+    value: Any = Field(serialization_alias="$const", 
validation_alias=AliasChoices(CONST_KEY))
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+def _value_discriminator(v: Any) -> str:
+    """
+    Route a raw value to a marker branch, or to ``literal``.
+
+    A reserved ``$``-marker is recognised only as the sole key of an object, so
+    a multi-key dict that merely contains ``$x`` is a plain literal dict.
+    """
+    if isinstance(v, dict) and len(v) == 1:
+        (key,) = v
+        if key in XCOM_KEYS:
+            return "xcom"
+        if key in TEMPLATE_KEYS:
+            return "template"
+        if key == CONST_KEY:
+            return "const"
+    # Already-parsed marker instances (revalidation) route to themselves.
+    if isinstance(v, XComRef):
+        return "xcom"
+    if isinstance(v, TemplateRef):
+        return "template"
+    if isinstance(v, ConstRef):
+        return "const"
+    return "literal"
+
+
+class _Literal(RootModel[Any]):
+    """
+    A literal value.
+
+    Containers recurse so nested markers are still resolved; scalars pass
+    through.
+    """
+
+    root: Any
+
+    @model_validator(mode="before")
+    @classmethod
+    def _recurse(cls, v: Any) -> Any:
+        if isinstance(v, dict):
+            return {k: _ValueAdapter.validate_python(item) for k, item in 
v.items()}
+        if isinstance(v, list):
+            return [_ValueAdapter.validate_python(item) for item in v]
+        return v
+
+
+Value = Annotated[
+    Annotated[XComRef, Tag("xcom")]
+    | Annotated[TemplateRef, Tag("template")]
+    | Annotated[ConstRef, Tag("const")]
+    | Annotated[_Literal, Tag("literal")],
+    Discriminator(_value_discriminator),
+]
+
+_ValueAdapter: TypeAdapter[Any] = TypeAdapter(Value)
+
+
+class _TaskBase(BaseModel):
+    """
+    Fields common to both ``use:`` and ``run:`` task kinds.
+
+    Unknown keys are BaseOperator arguments (pass-through).
+    """
+
+    id_: str = Field(alias="id")
+    needs: list[str] = Field(default_factory=list)
+    extends: list[str] = Field(default_factory=list)
+    with_: dict[str, Value] = Field(default_factory=dict, alias="with")
+
+    model_config = ConfigDict(extra="allow", populate_by_name=True)
+
+    @model_validator(mode="before")
+    @classmethod
+    def _check_task(cls, data: Any) -> Any:
+        if isinstance(data, dict):
+            for k in TASK_ASSET_IO_KEYS:
+                if k in data:
+                    raise ValueError(f"{k!r} not yet implemented")
+            for k in TASK_CALLBACK_KEYS:
+                if k in data:
+                    raise ValueError(f"{k!r} (callback) not yet implemented")
+            bodies = [k for k in ("uses", "run") if k in data]
+            if len(bodies) != 1:
+                raise ValueError(
+                    f"task {data.get('id')!r} must have exactly one of 'uses' 
or 'run'"
+                    + (f"; got {bodies}" if bodies else "")
+                )
+        return data
+
+
+class OperatorTask(_TaskBase):
+    """
+    A ready-made operator.
+
+    This should have ``uses`` (import path) + ``with``.
+    """
+
+    uses: str
+
+
+class CodeTask(_TaskBase):
+    """
+    A task running custom code.
+
+    This should have ``run`` holding function arguments, and optionally
+    ``queue`` to route it.
+    """
+
+    run: dict[str, Value] = Field(default_factory=dict)
+
+
+def _task_discriminator(v: Any) -> str:
+    if isinstance(v, OperatorTask):
+        return "operator"
+    if isinstance(v, CodeTask):
+        return "code"
+    if isinstance(v, dict) and "uses" in v:
+        return "operator"
+    return "code"
+
+
+Task = Annotated[
+    Annotated[OperatorTask, Tag("operator")] | Annotated[CodeTask, 
Tag("code")],
+    Discriminator(_task_discriminator),
+]
+
+
+class TaskTemplate(BaseModel):
+    """
+    A reusable, partial task fragment merged into any task that ``extends`` it.
+
+    This holds arbitrary keys as-is. All validation is done on the merged task.
+    """
+
+    model_config = ConfigDict(extra="allow")
+
+
+class TimetableSchedule(BaseModel):
+    """A constructed timetable."""
+
+    uses: str
+    with_: dict[str, Value] = Field(default_factory=dict, alias="with")
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+Schedule = str | None | TimetableSchedule
+
+
+def _merge_task(base: dict, over: dict) -> dict:
+    """
+    Shallow merge *over* onto *base*.
+
+    Note that ``with`` and ``run`` need to merge one level deep. All other
+    fields simply replace.
+    """
+    out = base.copy()
+    for k, v in over.items():
+        if k in ("with", "run") and isinstance(v, dict) and 
isinstance(out.get(k), dict):
+            merged = dict(out[k])
+            merged.update(v)
+            out[k] = merged
+        else:
+            out[k] = v
+    return out
+
+
+def _expand_template(name: str, templates: dict[str, Any], _seen: tuple = ()) 
-> dict:
+    if name in _seen:
+        raise ValueError(f"template cycle via {name!r}")
+    if name not in templates:
+        raise ValueError(f"unknown template {name!r}")
+    spec = templates[name]
+    if not isinstance(spec, dict):
+        raise ValueError(f"template {name!r} must be a mapping")
+    acc: dict = {}
+    for parent in spec.get("extends", []) or []:
+        acc = _merge_task(acc, _expand_template(parent, templates, _seen + 
(name,)))
+    own = {k: v for k, v in spec.items() if k != "extends"}
+    return _merge_task(acc, own)
+
+
+def _iter_xcom_task_ids(value: Any):
+    """Yield task ids referenced by any XComRef inside a resolved value."""
+    if isinstance(value, XComRef):
+        yield value.target if isinstance(value.target, str) else 
value.target.task
+    elif isinstance(value, _Literal):
+        yield from _iter_xcom_task_ids(value.root)
+    elif isinstance(value, dict):
+        for v in value.values():
+            yield from _iter_xcom_task_ids(v)
+    elif isinstance(value, list):
+        for v in value:
+            yield from _iter_xcom_task_ids(v)
+
+
+class DagDocument(BaseModel):
+    """
+    A Dag represented by one YAML document.
+
+    Extra top-level keys are Dag-level arguments.
+    """
+
+    schema_: str = Field(alias="$schema")
+    dag_id: str
+    schedule: Schedule = None
+    templates: dict[str, TaskTemplate] = Field(default_factory=dict)
+    tasks: list[Task] = Field(default_factory=list)
+
+    model_config = ConfigDict(extra="allow", populate_by_name=True)
+
+    @model_validator(mode="before")
+    @classmethod
+    def _reject_deferred(cls, data: Any) -> Any:
+        if isinstance(data, dict):
+            if "default_args" in data:
+                raise ValueError("default_args is not supported; use 
'templates' with 'extends' instead")
+            for k in DAG_CALLBACK_KEYS:
+                if k in data:
+                    raise ValueError(f"Dag-level {k!r} is not implemented yet")
+        return data
+
+    @model_validator(mode="before")
+    @classmethod
+    def _resolve_extends(cls, data: Any) -> Any:
+        """Merge in each task's ``extends`` templates before task 
validation."""
+        if not isinstance(data, dict):
+            return data
+        templates = data.get("templates")
+        if templates is None:
+            templates = {}
+        tasks = data.get("tasks")
+        # Only merge when the shapes are as expected. A malformed input is left
+        # untouched, so normal field validation reports it as a 
ValidationError.
+        if not isinstance(templates, dict) or not isinstance(tasks, list):
+            return data
+        resolved = []
+        for task in tasks:
+            extends = task.get("extends") if isinstance(task, dict) else None
+            if isinstance(extends, list) and extends:
+                merged: dict = {}
+                for name in extends:
+                    merged = _merge_task(merged, _expand_template(name, 
templates))
+                resolved.append(_merge_task(merged, {k: v for k, v in 
task.items() if k != "extends"}))
+            else:
+                resolved.append(task)
+        return {**data, "tasks": resolved}
+
+    @model_validator(mode="after")
+    def _check_refs(self):
+        """Check ``needs`` and XCom targets are present in this Dag."""
+        ids = {t.id_ for t in self.tasks}

Review Comment:
   Duplicate task ids get through: `tasks: [{id: a, run: {}}, {id: a, uses: 
x.Y}]` parses into two tasks, and `needs: [a]` then resolves against both. The 
Dag constructor rejects this eventually, but this validator already owns 
cross-task structure, and catching it here gives an error that points at the 
YAML. A count over the ids next to this set would be enough.



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/models.py:
##########
@@ -0,0 +1,362 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+Pydantic models for the YAML/JSON DAG format.
+
+These are the single source of truth for the format; they validate an authored
+document (rich Pydantic errors) and back the published JSON Schema. Only the
+format skeleton and value grammar only is described; operator ``with:`` args,
+Dag- and task-level attributes pass through untyped because the Python
+constructors validate those at import time.
+"""
+
+from __future__ import annotations
+
+import itertools
+from typing import Annotated, Any
+
+from pydantic import (
+    AliasChoices,
+    BaseModel,
+    ConfigDict,
+    Discriminator,
+    Field,
+    RootModel,
+    Tag,
+    TypeAdapter,
+    model_validator,
+)
+
+XCOM_KEYS = ("$x", "$xcom")
+TEMPLATE_KEYS = ("$t", "$template")
+CONST_KEY = "$const"
+
+# Keys deferred to a later edition...
+DAG_CALLBACK_KEYS = {"on_success_callback", "on_failure_callback", 
"sla_miss_callback"}

Review Comment:
   SLA was removed in Airflow 3 (`DAG.sla_miss_callback` is typed `None`, and 
`BaseOperator.sla` is marked deprecated), so "not implemented yet" suggests it 
is coming back. Rejecting `sla_miss_callback` with a "removed in Airflow 3, use 
deadline alerts" message, in this set and in the task one, would be more 
accurate.



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/models.py:
##########
@@ -0,0 +1,362 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+Pydantic models for the YAML/JSON DAG format.
+
+These are the single source of truth for the format; they validate an authored
+document (rich Pydantic errors) and back the published JSON Schema. Only the
+format skeleton and value grammar only is described; operator ``with:`` args,
+Dag- and task-level attributes pass through untyped because the Python
+constructors validate those at import time.
+"""
+
+from __future__ import annotations
+
+import itertools
+from typing import Annotated, Any
+
+from pydantic import (
+    AliasChoices,
+    BaseModel,
+    ConfigDict,
+    Discriminator,
+    Field,
+    RootModel,
+    Tag,
+    TypeAdapter,
+    model_validator,
+)
+
+XCOM_KEYS = ("$x", "$xcom")
+TEMPLATE_KEYS = ("$t", "$template")
+CONST_KEY = "$const"
+
+# Keys deferred to a later edition...
+DAG_CALLBACK_KEYS = {"on_success_callback", "on_failure_callback", 
"sla_miss_callback"}
+TASK_CALLBACK_KEYS = {
+    "on_success_callback",
+    "on_failure_callback",
+    "on_retry_callback",
+    "on_execute_callback",
+    "on_skipped_callback",
+    "sla_miss_callback",
+}
+TASK_ASSET_IO_KEYS = {"inlets", "outlets"}
+
+# Structural keys the format owns (everything else is pass-through to Python).
+_DAG_STRUCTURAL = {"$schema", "dag_id", "schedule", "templates", "tasks"}
+
+
+class XComTarget(BaseModel):
+    """The object form of an XCom reference."""
+
+    task: str
+    key: str | None = None
+
+    model_config = ConfigDict(extra="forbid")
+
+
+class XComRef(BaseModel):
+    """An upstream task's XCom output."""
+
+    target: str | XComTarget = Field(
+        serialization_alias="$x",
+        validation_alias=AliasChoices(*XCOM_KEYS),
+    )
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+class TemplateRef(BaseModel):
+    """A Jinja template."""
+
+    source: str = Field(serialization_alias="$t", 
validation_alias=AliasChoices(*TEMPLATE_KEYS))
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+class ConstRef(BaseModel):
+    """Force a literal; the value is taken verbatim."""
+
+    value: Any = Field(serialization_alias="$const", 
validation_alias=AliasChoices(CONST_KEY))
+
+    model_config = ConfigDict(extra="forbid", populate_by_name=True)
+
+
+def _value_discriminator(v: Any) -> str:
+    """
+    Route a raw value to a marker branch, or to ``literal``.
+
+    A reserved ``$``-marker is recognised only as the sole key of an object, so
+    a multi-key dict that merely contains ``$x`` is a plain literal dict.
+    """
+    if isinstance(v, dict) and len(v) == 1:

Review Comment:
   Since `$const` already exists to escape a literal `$` key, should an unknown 
`$`-prefixed sole key be rejected? Today a typo like `{$xcomm: up}` or 
`{$tempalte: ...}` silently becomes a literal dict passed to the operator. The 
XCom edge and the unknown-task check both disappear without an error. The same 
goes for `{$x: a, key: k}`, which reads like a keyed XCom reference but is a 
literal.



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/migrator.py:
##########
@@ -0,0 +1,132 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+`$schema` version resolution and forward migration to the head shape.
+
+A document pinned to an older `$schema` version is migrated up to the current 
head shape
+*before* validation, by walking the Cadwyn version bundle and applying each 
`VersionChange`'s
+forward request converters 
(`@convert_request_to_next_version_for(DagDocument)`), oldest to
+newest. This mirrors the exec-API / supervisor-schema `SchemaVersionMigrator`, 
trimmed to the
+one direction a document format needs (author's version -> head), with the 
bundle and its
+version list encapsulated in the migrator rather than exposed as module-level 
helpers.
+
+With a single published version there is nothing to migrate; the machinery is 
here so a future
+version bump only has to add a `VersionChange` carrying a converter — no 
parser changes.
+"""
+
+from __future__ import annotations
+
+import copy
+import functools
+import re
+import warnings
+from typing import TYPE_CHECKING, Any
+
+import attrs
+
+from airflow.sdk.importers.yaml_importer.models import DagDocument
+from airflow.sdk.importers.yaml_importer.versions import get_bundle
+
+if TYPE_CHECKING:
+    from cadwyn import VersionBundle
+
+_SCHEMA_URL = "https://airflow.apache.org/schemas/dag/{date}.json";

Review Comment:
   Nothing publishes `https://airflow.apache.org/schemas/dag/2026-10-30.json` 
yet: `publish-schemas-to-s3` only takes the Execution API and supervisor 
schemas. There is also no counterpart to `check-supervisor-schemas-versions`, 
so a later models.py change would rewrite the 2026-10-30 schema in place. Is 
that planned as a follow-up? If so, it's worth a line in the description, since 
editors that fetch `$schema` will get a 404 until then.



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/migrator.py:
##########
@@ -0,0 +1,132 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+`$schema` version resolution and forward migration to the head shape.
+
+A document pinned to an older `$schema` version is migrated up to the current 
head shape
+*before* validation, by walking the Cadwyn version bundle and applying each 
`VersionChange`'s
+forward request converters 
(`@convert_request_to_next_version_for(DagDocument)`), oldest to
+newest. This mirrors the exec-API / supervisor-schema `SchemaVersionMigrator`, 
trimmed to the
+one direction a document format needs (author's version -> head), with the 
bundle and its
+version list encapsulated in the migrator rather than exposed as module-level 
helpers.
+
+With a single published version there is nothing to migrate; the machinery is 
here so a future
+version bump only has to add a `VersionChange` carrying a converter — no 
parser changes.
+"""
+
+from __future__ import annotations
+
+import copy
+import functools
+import re
+import warnings
+from typing import TYPE_CHECKING, Any
+
+import attrs
+
+from airflow.sdk.importers.yaml_importer.models import DagDocument
+from airflow.sdk.importers.yaml_importer.versions import get_bundle
+
+if TYPE_CHECKING:
+    from cadwyn import VersionBundle
+
+_SCHEMA_URL = "https://airflow.apache.org/schemas/dag/{date}.json";
+_VERSION_RE = re.compile(r"\d{4}-\d{2}-\d{2}")
+
+
+def schema_url(date: str) -> str:
+    """Build the canonical ``$schema`` URL for a published version *date*."""
+    return _SCHEMA_URL.format(date=date)
+
+
+def version_from_schema(schema: str) -> str | None:
+    """Read the version (``YYYY-MM-DD``) token out of a ``$schema`` URL, 
ignoring the host."""
+    match = _VERSION_RE.search(schema or "")
+    return match.group(0) if match else None
+
+
[email protected]
+class _RequestInfo:
+    """
+    Duck-type stand-in for Cadwyn's ``RequestInfo``.
+
+    ``cadwyn.structure.data._AlterDataInstruction.__call__`` only reads
+    and writes ``info.body``; the by-schema transformers we drive never
+    touch FastAPI's Request/Response. Passing this minimal object lets
+    us run cadwyn's migrations from a pure in-process code path with no
+    HTTP stack.
+    """
+
+    body: dict[str, Any]
+
+
[email protected]
+class DagDocumentMigrator:
+    """YAML Dag document migrator; pins each document to its ``$schema`` 
version."""
+
+    _bundle: VersionBundle
+
+    def resolve_and_migrate(self, body: dict[str, Any], *, source: str) -> 
dict[str, Any]:
+        """
+        Resolve *body*'s ``$schema`` version and migrate it to the head shape.
+
+        The version token is read from the ``$schema`` URL (the host is 
ignored). It is an
+        exact pin: a known version uses its own ruleset; an unknown one (e.g. 
newer than this
+        importer) resolves to the latest ruleset with a warning, never a hard 
failure.
+
+        :return: The migrated result. If *body* is already at head, it is
+            returned as-is; otherwise a migrated copy is returned.
+        """
+        known = [v.value for v in self._bundle.versions]  # newest-first
+        date = version_from_schema(body["$schema"])
+        if date in known:
+            source_version = date
+        else:
+            source_version = known[0]  # unknown version -> latest ruleset we 
have

Review Comment:
   This fallback catches more than versions newer than the importer. It also 
takes dates older than the oldest version, dates between two versions, and a 
`$schema` with no date at all (`$schema: hello` warns about version `None` and 
parses). Once a second version exists, a document with an unpublished date 
older than head gets the head ruleset and skips the converters it needs. 
Because of `extra="allow"`, its old keys then pass through as kwargs. Would it 
be safer to fall back to head only for dates newer than head, and raise for a 
missing date or an older date that was never published?



##########
task-sdk/src/airflow/sdk/importers/yaml_importer/__init__.py:
##########
@@ -0,0 +1,38 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""
+Native YAML/JSON DAG format for Airflow (`$schema` version 2026-10-30).
+
+Public API:

Review Comment:
   This calls itself the public API, but the importer base it belongs to is 
marked `|experimental|` (`AbstractDagImporter`, base.py:233). The format is 
still being settled on the dev list and several keys raise "not yet 
implemented", so should this package carry the same marker, keeping the 
2026-10-30 shape free to change?



##########
task-sdk/tests/task_sdk/importers/yaml_importer/test_parser.py:
##########
@@ -0,0 +1,328 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""Unit tests for the YAML DAG format parser 
(airflow.sdk.importers.yaml_importer)."""
+
+from __future__ import annotations
+
+import io
+import textwrap
+
+import pytest
+
+from airflow.sdk.importers.yaml_importer import YamlDagParseError, 
parse_documents
+from airflow.sdk.importers.yaml_importer.models import (
+    CodeTask,
+    ConstRef,
+    DagDocument,
+    OperatorTask,
+    TemplateRef,
+    TimetableSchedule,
+    XComRef,
+    XComTarget,
+    _Literal,
+)
+
+HEAD = "2026-10-30"
+HEAD_SCHEMA = f"https://airflow.apache.org/schemas/dag/{HEAD}.json";
+
+
+def _one(body: str):
+    return next(iter(parse_documents(f"$schema: {HEAD_SCHEMA}\ndag_id: 
d\n{textwrap.dedent(body)}")))
+
+
+def _task(task_yaml: str):
+    return _one("tasks:\n" + textwrap.indent(textwrap.dedent(task_yaml), "  
")).tasks[0]
+
+
+# --------------------------------------------------------------------------- 
value grammar
+def test_literal_by_default():
+    t = _task("- {id: t, run: {region: 'eu', n: 5, flag: true}}")
+    assert isinstance(t.run["region"], _Literal)
+    assert t.run["region"].root == "eu"
+    assert t.run["n"].root == 5
+    assert t.run["flag"].root is True
+
+
+def test_xcom_short_long_and_keyed():
+    doc = _one(
+        "tasks:\n  - {id: up, run: {}}\n  - {id: t, run: {a: {$x: up}, b: 
{$xcom: up}, c: {$x: {task: up, key: rows}}}}"
+    )
+    t = {x.id_: x for x in doc.tasks}["t"]
+    assert isinstance(t.run["a"], XComRef)
+    assert t.run["a"].target == "up"
+    assert isinstance(t.run["b"], XComRef)
+    assert t.run["b"].target == "up"
+    assert isinstance(t.run["c"].target, XComTarget)
+    assert t.run["c"].target.key == "rows"
+
+
+def test_template_marker():
+    t = _task("- {id: t, uses: X, with: {k: {$t: 'v-{{ ds }}'}, k2: 
{$template: 'x'}}}")
+    assert isinstance(t.with_["k"], TemplateRef)
+    assert "{{ ds }}" in t.with_["k"].source
+    assert isinstance(t.with_["k2"], TemplateRef)
+
+
+def test_const_is_verbatim_not_recursed():
+    t = _task("- {id: t, run: {a: {$const: {$x: not-a-ref}}}}")
+    assert isinstance(t.run["a"], ConstRef)
+    assert t.run["a"].value == {"$x": "not-a-ref"}  # inner marker NOT 
interpreted
+
+
+def test_marker_only_as_sole_key():
+    t = _task("- {id: t, run: {a: {$x: up, extra: 1}}}")  # two keys -> 
literal dict
+    assert isinstance(t.run["a"], _Literal)
+    assert set(t.run["a"].root) == {"$x", "extra"}
+
+
+def test_nested_marker_inside_literal_resolved():
+    t = _task("- {id: t, uses: X, with: {cfg: {url: {$t: '{{ ds }}'}, n: 1}}}")
+    inner = t.with_["cfg"].root
+    assert isinstance(inner["url"], TemplateRef)
+    assert inner["n"].root == 1
+
+
+# --------------------------------------------------------------------------- 
task bodies
+def test_operator_vs_code_discrimination():
+    assert isinstance(_task("- {id: t, uses: a.b.C, with: {x: 1}}"), 
OperatorTask)
+    assert isinstance(_task("- {id: t, run: {}}"), CodeTask)
+
+
+def test_no_body_is_error():
+    with pytest.raises(YamlDagParseError, match="exactly one"):
+        _task("- {id: t}")
+
+
+def test_both_bodies_is_error():
+    with pytest.raises(YamlDagParseError, match="exactly one"):
+        _task("- {id: t, uses: X, run: {}}")
+
+
+# --------------------------------------------------------------------------- 
deferred keys
+def test_default_args_rejected():
+    with pytest.raises(YamlDagParseError, match="default_args"):
+        _one("default_args: {retries: 1}\ntasks: []")
+
+
+def test_callbacks_rejected():
+    with pytest.raises(YamlDagParseError, match="callback"):
+        _one("on_failure_callback: cb\ntasks: []")
+    with pytest.raises(YamlDagParseError, match="callback"):
+        _task("- {id: t, run: {}, on_success_callback: cb}")
+
+
+def test_inlets_outlets_rejected():
+    with pytest.raises(YamlDagParseError, match="not yet implemented"):
+        _task("- {id: t, run: {}, outlets: [a]}")
+
+
+# --------------------------------------------------------------------------- 
templates / extends
+def test_extends_merges_shallow_task_wins():
+    doc = _one(
+        """
+        templates:
+          retryable: {retries: 2, retry_delay: '5m'}
+          load:
+            extends: [retryable]
+            uses: a.b.S3ToRedshiftOperator
+            with: {schema: public, conn_id: warehouse}
+        tasks:
+          - id: t
+            extends: [load]
+            with: {table: sales, conn_id: override}
+        """
+    )
+    t = doc.tasks[0]
+    assert isinstance(t, OperatorTask)
+    assert t.uses.endswith("S3ToRedshiftOperator")
+    assert t.__pydantic_extra__["retries"] == 2  # from retryable
+    assert isinstance(t.with_["schema"], _Literal)  # inherited
+    assert t.with_["table"].root == "sales"  # own
+    assert t.with_["conn_id"].root == "override"  # one-level with-merge, task 
wins
+
+
+def test_unknown_template_errors():
+    with pytest.raises(YamlDagParseError, match="unknown template"):
+        _one("tasks:\n  - {id: t, extends: [ghost], run: {}}")
+
+
+def test_template_cycle_errors():
+    with pytest.raises(YamlDagParseError, match="cycle"):
+        _one("templates: {a: {extends: [b]}, b: {extends: [a]}}\ntasks: [{id: 
t, extends: [a], run: {}}]")
+
+
+# --------------------------------------------------------------------------- 
references / edges
+def test_unknown_needs_rejected():
+    with pytest.raises(YamlDagParseError, match="unknown task"):
+        _one("tasks:\n  - {id: t, run: {}, needs: [ghost]}")
+
+
+def test_unknown_xcom_target_rejected():
+    with pytest.raises(YamlDagParseError, match="unknown task"):
+        _one("tasks:\n  - {id: t, run: {e: {$x: ghost}}}")
+
+
+# --------------------------------------------------------------------------- 
schedule / dag attrs
+def test_schedule_forms():
+    assert _one("schedule: '0 3 * * *'\ntasks: []").schedule == "0 3 * * *"
+    assert _one("schedule: null\ntasks: []").schedule is None
+    tt = _one("schedule: {uses: a.b.MyTimetable, with: {n: 1}}\ntasks: 
[]").schedule
+    assert isinstance(tt, TimetableSchedule)
+    assert tt.uses.endswith("MyTimetable")
+
+
+def test_dag_attributes_pass_through():
+    doc = _one("catchup: false\nmax_active_runs: 3\ntags: [etl]\ntasks: []")
+    assert doc.dag_attributes == {"catchup": False, "max_active_runs": 3, 
"tags": ["etl"]}
+
+
+# --------------------------------------------------------------------------- 
$schema
+def test_schema_required():
+    with pytest.raises(YamlDagParseError, match=r"\$schema"):
+        list(parse_documents("dag_id: d\ntasks: []"))
+
+
+def test_unknown_date_warns_and_falls_back():
+    url = "https://airflow.apache.org/schemas/dag/2099-01-01.json";
+    with pytest.warns(UserWarning, match="not a known version"):
+        doc = next(iter(parse_documents(f"$schema: {url}\ndag_id: d\ntasks: 
[]")))
+    assert doc.dag_id == "d"
+
+
+# --------------------------------------------------------------------------- 
multi-document
+def test_accepts_a_file_like_stream():
+    stream = io.StringIO(f"$schema: {HEAD_SCHEMA}\ndag_id: d\ntasks: [{{id: t, 
run: {{}}}}]")
+    docs = list(parse_documents(stream))
+    assert [d.dag_id for d in docs] == ["d"]
+
+
+def test_parse_is_lazy_errors_surface_on_iteration():
+    # Building the iterator does not parse; the error only surfaces when 
consumed.
+    gen = parse_documents("dag_id: d\ntasks: []")  # missing $schema
+    with pytest.raises(YamlDagParseError, match=r"\$schema"):
+        next(iter(gen))
+
+
+def test_multiple_documents():
+    docs = parse_documents(
+        f"$schema: {HEAD_SCHEMA}\ndag_id: a\ntasks: []\n---\n$schema: 
{HEAD_SCHEMA}\ndag_id: b\ntasks: []\n"
+    )
+    assert [d.dag_id for d in docs] == ["a", "b"]
+
+
+# --------------------------------------------------------------------------- 
schema
+def test_model_json_schema_describes_markers():
+    # The models are the schema source of truth (the prek dump script 
decorates + snapshots
+    # them). Assert the generated schema shape here; the snapshot itself is 
enforced by the
+    # generate-yaml-importer-schema-snapshot prek hook.
+    sch = DagDocument.model_json_schema(by_alias=True)
+    assert set(sch["properties"]) >= {"$schema", "dag_id", "schedule", 
"tasks", "templates"}
+    assert {"XComRef", "TemplateRef", "ConstRef", "OperatorTask", "CodeTask"} 
<= set(sch["$defs"])
+
+
+# --------------------------------------------------------------------------- 
integration
+def test_parses_example_fixture():
+    import pathlib
+
+    fixture = pathlib.Path(__file__).parent / "example_dag.yaml"
+    with fixture.open(encoding="utf-8") as fh:
+        docs = list(parse_documents(fh, source=str(fixture)))
+    assert [d.dag_id for d in docs] == ["retail_daily_sales", 
"events_pipeline"]
+
+    d1 = {t.id_: t for t in docs[0].tasks}
+    load = d1["load_sales"]
+    assert isinstance(load, OperatorTask)
+    assert load.uses.endswith("S3ToRedshiftOperator")
+    assert load.__pydantic_extra__["retries"] == 2  # via redshift_load -> 
retryable
+    assert isinstance(load.with_["schema"], TemplateRef)  # inherited $t
+    assert isinstance(load.with_["s3_key"], TemplateRef)  # own $t
+    assert load.needs == ["wait_for_export"]
+    assert docs[0].dag_attributes["catchup"] is False
+
+    d2 = {t.id_: t for t in docs[1].tasks}
+    assert isinstance(d2["extract"], CodeTask)
+    assert d2["extract"].run == {}
+    assert d2["extract"].queue == "extract-workers"
+    assert isinstance(d2["transform"].run["events"], XComRef)  # {$x: extract} 
edge
+    assert d2["load"].needs == ["transform"]
+
+
+# --------------------------------------------------------------------------- 
$schema version migration
+# A synthetic two-version bundle with a real forward converter, to exercise 
the migration
+# machinery (the real bundle has a single version, so nothing migrates there).
+from cadwyn import (  # noqa: E402

Review Comment:
   Nothing forces these imports to sit mid-module: `migrator` is already loaded 
through the import at the top, and cadwyn imports fine at load time. `import 
pathlib` on line 239 is the same. Could they move to the top so the `noqa` can 
go? Moving the migrator tests into their own `test_migrator.py`, as the sibling 
schema package does, would also take care of it.



##########
task-sdk/tests/task_sdk/importers/yaml_importer/test_parser.py:
##########
@@ -0,0 +1,328 @@
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""Unit tests for the YAML DAG format parser 
(airflow.sdk.importers.yaml_importer)."""
+
+from __future__ import annotations
+
+import io
+import textwrap
+
+import pytest
+
+from airflow.sdk.importers.yaml_importer import YamlDagParseError, 
parse_documents
+from airflow.sdk.importers.yaml_importer.models import (
+    CodeTask,
+    ConstRef,
+    DagDocument,
+    OperatorTask,
+    TemplateRef,
+    TimetableSchedule,
+    XComRef,
+    XComTarget,
+    _Literal,
+)
+
+HEAD = "2026-10-30"
+HEAD_SCHEMA = f"https://airflow.apache.org/schemas/dag/{HEAD}.json";
+
+
+def _one(body: str):
+    return next(iter(parse_documents(f"$schema: {HEAD_SCHEMA}\ndag_id: 
d\n{textwrap.dedent(body)}")))
+
+
+def _task(task_yaml: str):
+    return _one("tasks:\n" + textwrap.indent(textwrap.dedent(task_yaml), "  
")).tasks[0]
+
+
+# --------------------------------------------------------------------------- 
value grammar
+def test_literal_by_default():
+    t = _task("- {id: t, run: {region: 'eu', n: 5, flag: true}}")
+    assert isinstance(t.run["region"], _Literal)
+    assert t.run["region"].root == "eu"
+    assert t.run["n"].root == 5
+    assert t.run["flag"].root is True
+
+
+def test_xcom_short_long_and_keyed():
+    doc = _one(
+        "tasks:\n  - {id: up, run: {}}\n  - {id: t, run: {a: {$x: up}, b: 
{$xcom: up}, c: {$x: {task: up, key: rows}}}}"
+    )
+    t = {x.id_: x for x in doc.tasks}["t"]
+    assert isinstance(t.run["a"], XComRef)
+    assert t.run["a"].target == "up"
+    assert isinstance(t.run["b"], XComRef)
+    assert t.run["b"].target == "up"
+    assert isinstance(t.run["c"].target, XComTarget)
+    assert t.run["c"].target.key == "rows"
+
+
+def test_template_marker():
+    t = _task("- {id: t, uses: X, with: {k: {$t: 'v-{{ ds }}'}, k2: 
{$template: 'x'}}}")
+    assert isinstance(t.with_["k"], TemplateRef)
+    assert "{{ ds }}" in t.with_["k"].source
+    assert isinstance(t.with_["k2"], TemplateRef)
+
+
+def test_const_is_verbatim_not_recursed():
+    t = _task("- {id: t, run: {a: {$const: {$x: not-a-ref}}}}")
+    assert isinstance(t.run["a"], ConstRef)
+    assert t.run["a"].value == {"$x": "not-a-ref"}  # inner marker NOT 
interpreted
+
+
+def test_marker_only_as_sole_key():
+    t = _task("- {id: t, run: {a: {$x: up, extra: 1}}}")  # two keys -> 
literal dict
+    assert isinstance(t.run["a"], _Literal)
+    assert set(t.run["a"].root) == {"$x", "extra"}
+
+
+def test_nested_marker_inside_literal_resolved():
+    t = _task("- {id: t, uses: X, with: {cfg: {url: {$t: '{{ ds }}'}, n: 1}}}")
+    inner = t.with_["cfg"].root
+    assert isinstance(inner["url"], TemplateRef)
+    assert inner["n"].root == 1
+
+
+# --------------------------------------------------------------------------- 
task bodies
+def test_operator_vs_code_discrimination():
+    assert isinstance(_task("- {id: t, uses: a.b.C, with: {x: 1}}"), 
OperatorTask)
+    assert isinstance(_task("- {id: t, run: {}}"), CodeTask)
+
+
+def test_no_body_is_error():
+    with pytest.raises(YamlDagParseError, match="exactly one"):
+        _task("- {id: t}")
+
+
+def test_both_bodies_is_error():
+    with pytest.raises(YamlDagParseError, match="exactly one"):
+        _task("- {id: t, uses: X, run: {}}")
+
+
+# --------------------------------------------------------------------------- 
deferred keys
+def test_default_args_rejected():
+    with pytest.raises(YamlDagParseError, match="default_args"):
+        _one("default_args: {retries: 1}\ntasks: []")
+
+
+def test_callbacks_rejected():
+    with pytest.raises(YamlDagParseError, match="callback"):
+        _one("on_failure_callback: cb\ntasks: []")
+    with pytest.raises(YamlDagParseError, match="callback"):
+        _task("- {id: t, run: {}, on_success_callback: cb}")
+
+
+def test_inlets_outlets_rejected():
+    with pytest.raises(YamlDagParseError, match="not yet implemented"):
+        _task("- {id: t, run: {}, outlets: [a]}")
+
+
+# --------------------------------------------------------------------------- 
templates / extends
+def test_extends_merges_shallow_task_wins():
+    doc = _one(
+        """
+        templates:
+          retryable: {retries: 2, retry_delay: '5m'}
+          load:
+            extends: [retryable]
+            uses: a.b.S3ToRedshiftOperator
+            with: {schema: public, conn_id: warehouse}
+        tasks:
+          - id: t
+            extends: [load]
+            with: {table: sales, conn_id: override}
+        """
+    )
+    t = doc.tasks[0]
+    assert isinstance(t, OperatorTask)
+    assert t.uses.endswith("S3ToRedshiftOperator")
+    assert t.__pydantic_extra__["retries"] == 2  # from retryable
+    assert isinstance(t.with_["schema"], _Literal)  # inherited
+    assert t.with_["table"].root == "sales"  # own
+    assert t.with_["conn_id"].root == "override"  # one-level with-merge, task 
wins
+
+
+def test_unknown_template_errors():
+    with pytest.raises(YamlDagParseError, match="unknown template"):
+        _one("tasks:\n  - {id: t, extends: [ghost], run: {}}")
+
+
+def test_template_cycle_errors():
+    with pytest.raises(YamlDagParseError, match="cycle"):
+        _one("templates: {a: {extends: [b]}, b: {extends: [a]}}\ntasks: [{id: 
t, extends: [a], run: {}}]")
+
+
+# --------------------------------------------------------------------------- 
references / edges
+def test_unknown_needs_rejected():
+    with pytest.raises(YamlDagParseError, match="unknown task"):
+        _one("tasks:\n  - {id: t, run: {}, needs: [ghost]}")
+
+
+def test_unknown_xcom_target_rejected():
+    with pytest.raises(YamlDagParseError, match="unknown task"):
+        _one("tasks:\n  - {id: t, run: {e: {$x: ghost}}}")
+
+
+# --------------------------------------------------------------------------- 
schedule / dag attrs
+def test_schedule_forms():
+    assert _one("schedule: '0 3 * * *'\ntasks: []").schedule == "0 3 * * *"
+    assert _one("schedule: null\ntasks: []").schedule is None
+    tt = _one("schedule: {uses: a.b.MyTimetable, with: {n: 1}}\ntasks: 
[]").schedule
+    assert isinstance(tt, TimetableSchedule)
+    assert tt.uses.endswith("MyTimetable")
+
+
+def test_dag_attributes_pass_through():
+    doc = _one("catchup: false\nmax_active_runs: 3\ntags: [etl]\ntasks: []")
+    assert doc.dag_attributes == {"catchup": False, "max_active_runs": 3, 
"tags": ["etl"]}
+
+
+# --------------------------------------------------------------------------- 
$schema
+def test_schema_required():
+    with pytest.raises(YamlDagParseError, match=r"\$schema"):
+        list(parse_documents("dag_id: d\ntasks: []"))
+
+
+def test_unknown_date_warns_and_falls_back():
+    url = "https://airflow.apache.org/schemas/dag/2099-01-01.json";
+    with pytest.warns(UserWarning, match="not a known version"):
+        doc = next(iter(parse_documents(f"$schema: {url}\ndag_id: d\ntasks: 
[]")))
+    assert doc.dag_id == "d"
+
+
+# --------------------------------------------------------------------------- 
multi-document
+def test_accepts_a_file_like_stream():
+    stream = io.StringIO(f"$schema: {HEAD_SCHEMA}\ndag_id: d\ntasks: [{{id: t, 
run: {{}}}}]")
+    docs = list(parse_documents(stream))
+    assert [d.dag_id for d in docs] == ["d"]
+
+
+def test_parse_is_lazy_errors_surface_on_iteration():
+    # Building the iterator does not parse; the error only surfaces when 
consumed.
+    gen = parse_documents("dag_id: d\ntasks: []")  # missing $schema
+    with pytest.raises(YamlDagParseError, match=r"\$schema"):
+        next(iter(gen))
+
+
+def test_multiple_documents():
+    docs = parse_documents(
+        f"$schema: {HEAD_SCHEMA}\ndag_id: a\ntasks: []\n---\n$schema: 
{HEAD_SCHEMA}\ndag_id: b\ntasks: []\n"
+    )
+    assert [d.dag_id for d in docs] == ["a", "b"]
+
+
+# --------------------------------------------------------------------------- 
schema
+def test_model_json_schema_describes_markers():
+    # The models are the schema source of truth (the prek dump script 
decorates + snapshots
+    # them). Assert the generated schema shape here; the snapshot itself is 
enforced by the
+    # generate-yaml-importer-schema-snapshot prek hook.
+    sch = DagDocument.model_json_schema(by_alias=True)
+    assert set(sch["properties"]) >= {"$schema", "dag_id", "schedule", 
"tasks", "templates"}
+    assert {"XComRef", "TemplateRef", "ConstRef", "OperatorTask", "CodeTask"} 
<= set(sch["$defs"])
+
+
+# --------------------------------------------------------------------------- 
integration
+def test_parses_example_fixture():
+    import pathlib
+
+    fixture = pathlib.Path(__file__).parent / "example_dag.yaml"
+    with fixture.open(encoding="utf-8") as fh:
+        docs = list(parse_documents(fh, source=str(fixture)))
+    assert [d.dag_id for d in docs] == ["retail_daily_sales", 
"events_pipeline"]
+
+    d1 = {t.id_: t for t in docs[0].tasks}
+    load = d1["load_sales"]
+    assert isinstance(load, OperatorTask)
+    assert load.uses.endswith("S3ToRedshiftOperator")
+    assert load.__pydantic_extra__["retries"] == 2  # via redshift_load -> 
retryable
+    assert isinstance(load.with_["schema"], TemplateRef)  # inherited $t
+    assert isinstance(load.with_["s3_key"], TemplateRef)  # own $t
+    assert load.needs == ["wait_for_export"]
+    assert docs[0].dag_attributes["catchup"] is False
+
+    d2 = {t.id_: t for t in docs[1].tasks}
+    assert isinstance(d2["extract"], CodeTask)
+    assert d2["extract"].run == {}
+    assert d2["extract"].queue == "extract-workers"
+    assert isinstance(d2["transform"].run["events"], XComRef)  # {$x: extract} 
edge
+    assert d2["load"].needs == ["transform"]
+
+
+# --------------------------------------------------------------------------- 
$schema version migration
+# A synthetic two-version bundle with a real forward converter, to exercise 
the migration
+# machinery (the real bundle has a single version, so nothing migrates there).
+from cadwyn import (  # noqa: E402
+    HeadVersion,
+    Version,
+    VersionBundle,
+    VersionChange,
+    convert_request_to_next_version_for,
+)
+
+from airflow.sdk.importers.yaml_importer import migrator  # noqa: E402
+
+
+class _RenameOwnerAttr(VersionChange):
+    "test-only: 2026-10-30 spelled the attribute `owner_old`; head renamed it 
to `owner`."
+
+    description = __doc__
+    instructions_to_migrate_to_previous_version = ()
+
+    @convert_request_to_next_version_for(DagDocument)  # type: ignore[arg-type]
+    def _move(request):
+        if "owner_old" in request.body:
+            request.body["owner"] = request.body.pop("owner_old")
+
+
+_SYNTH_BUNDLE = VersionBundle(HeadVersion(), Version("2027-06-01", 
_RenameOwnerAttr), Version("2026-10-30"))

Review Comment:
   With two versions the only converter always sits on head, so the `<=` skip 
in `_migrate_to_head` is never exercised: changing it to `<` keeps all 30 tests 
green. Adding a third version with a non-idempotent converter and pinning the 
middle one would cover that boundary, and could also assert that the input body 
isn't mutated. `execution_time/schema/test_migrator.py` uses four versions for 
the same reason.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to