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


##########
airflow-core/src/airflow/serialization/serialized_objects.py:
##########
@@ -2131,6 +2134,70 @@ def from_dict(cls, serialized_obj: dict) -> 
SerializedDAG:
         # Pass client_defaults directly to deserialize_dag
         return cls.deserialize_dag(serialized_obj["dag"], client_defaults)
 
+    @classmethod
+    def fill_config_defaults(cls, serialized_obj: dict[str, Any]) -> None:
+        """
+        Fill in the Dag settings a serialized Dag leaves unset from the 
Airflow config, as a Python Dag does.
+
+        A Lang-SDK runtime cannot read the Airflow config, so it leaves 
``max_active_tasks``,
+        ``max_active_runs``, ``max_consecutive_failed_dag_runs``, ``catchup`` 
and
+        ``disable_bundle_versioning`` out unless the Dag sets them. A value 
the Dag sets is kept.
+        *serialized_obj* is changed in place.
+        """
+        dag = serialized_obj.get("dag")
+        if not isinstance(dag, dict):
+            # validate_serialized_dag rejects it.
+            return
+        for key, get, section, option in (
+            ("max_active_tasks", conf.getint, "core", 
"max_active_tasks_per_dag"),
+            ("max_active_runs", conf.getint, "core", 
"max_active_runs_per_dag"),
+            (
+                "max_consecutive_failed_dag_runs",
+                conf.getint,
+                "core",
+                "max_consecutive_failed_dag_runs_per_dag",
+            ),
+            ("catchup", conf.getboolean, "scheduler", "catchup_by_default"),
+            ("disable_bundle_versioning", conf.getboolean, "dag_processor", 
"disable_bundle_versioning"),
+        ):
+            if key not in dag:
+                dag[key] = get(section, option)
+
+    @classmethod
+    def validate_serialized_dag(cls, serialized_obj: dict[str, Any]) -> None:
+        """
+        Check that a serialized Dag, such as one a Lang-SDK runtime produced, 
can be stored and loaded.
+
+        It must match the JSON schema, deserialize, and have no cycle in its 
task graph.
+        *serialized_obj* is not changed.
+
+        :raises DeserializationError: if it does not.
+        """
+        from jsonschema import ValidationError
+
+        dag = serialized_obj.get("dag")
+        dag_id = dag.get("dag_id") if isinstance(dag, dict) else None
+        try:
+            cls.validate_schema(serialized_obj)

Review Comment:
   The `tasks` definition in schema.json puts `additionalProperties` on an 
array, which JSON Schema ignores, so this step never checks a task entry. A 
`{}` in `tasks` passes all three checks here (`from_dict` and the cycle walk 
both skip non-operator entries), and then `LazyDeserializedDAG.owner` raises 
`KeyError: '__var'` when the Dag processor writes the DagModel row 
(collection.py:751), outside the import-error path this feeds. Could this 
require every entry to be an operator with a `task_id`, and reject a repeated 
`task_id` while it's there? Today `from_dict` and the dict below both keep the 
last one silently.



##########
airflow-core/src/airflow/serialization/serialized_objects.py:
##########
@@ -2131,6 +2134,70 @@ def from_dict(cls, serialized_obj: dict) -> 
SerializedDAG:
         # Pass client_defaults directly to deserialize_dag
         return cls.deserialize_dag(serialized_obj["dag"], client_defaults)
 
+    @classmethod
+    def fill_config_defaults(cls, serialized_obj: dict[str, Any]) -> None:
+        """
+        Fill in the Dag settings a serialized Dag leaves unset from the 
Airflow config, as a Python Dag does.
+
+        A Lang-SDK runtime cannot read the Airflow config, so it leaves 
``max_active_tasks``,
+        ``max_active_runs``, ``max_consecutive_failed_dag_runs``, ``catchup`` 
and
+        ``disable_bundle_versioning`` out unless the Dag sets them. A value 
the Dag sets is kept.
+        *serialized_obj* is changed in place.
+        """
+        dag = serialized_obj.get("dag")
+        if not isinstance(dag, dict):
+            # validate_serialized_dag rejects it.
+            return
+        for key, get, section, option in (
+            ("max_active_tasks", conf.getint, "core", 
"max_active_tasks_per_dag"),

Review Comment:
   `scripts/ci/lang_sdk_serialization/serialize_python.py` keeps its own 
`CONFIG_BACKED_DAG_FIELDS` with these five entries, and its `receive()` repeats 
the schema check and `from_dict` that `validate_serialized_dag` now does. Could 
`receive()` call `fill_config_defaults` and `validate_serialized_dag` instead, 
so the conformance job checks the code the Dag processor runs and the list only 
lives here?



##########
airflow-core/src/airflow/serialization/serialized_objects.py:
##########
@@ -2131,6 +2134,70 @@ def from_dict(cls, serialized_obj: dict) -> 
SerializedDAG:
         # Pass client_defaults directly to deserialize_dag
         return cls.deserialize_dag(serialized_obj["dag"], client_defaults)
 
+    @classmethod
+    def fill_config_defaults(cls, serialized_obj: dict[str, Any]) -> None:

Review Comment:
   This fills the Dag-level settings, but a runtime can't read the task-level 
config either. A Python task gets `[core] default_task_retries`, `[operators] 
default_queue` and the other operator defaults through the payload's 
`client_defaults`, which a Lang-SDK payload never carries, so with 
`default_task_retries = 3` the same task loads with 3 retries from Python and 0 
from the TS runtime. Is that a planned follow-up? Filling `client_defaults` 
here would also need each runtime to keep a task field the author set 
explicitly, even when it equals the default.



##########
airflow-core/src/airflow/serialization/serialized_objects.py:
##########
@@ -2131,6 +2134,70 @@ def from_dict(cls, serialized_obj: dict) -> 
SerializedDAG:
         # Pass client_defaults directly to deserialize_dag
         return cls.deserialize_dag(serialized_obj["dag"], client_defaults)
 
+    @classmethod
+    def fill_config_defaults(cls, serialized_obj: dict[str, Any]) -> None:
+        """
+        Fill in the Dag settings a serialized Dag leaves unset from the 
Airflow config, as a Python Dag does.
+
+        A Lang-SDK runtime cannot read the Airflow config, so it leaves 
``max_active_tasks``,
+        ``max_active_runs``, ``max_consecutive_failed_dag_runs``, ``catchup`` 
and
+        ``disable_bundle_versioning`` out unless the Dag sets them. A value 
the Dag sets is kept.
+        *serialized_obj* is changed in place.
+        """
+        dag = serialized_obj.get("dag")
+        if not isinstance(dag, dict):
+            # validate_serialized_dag rejects it.
+            return
+        for key, get, section, option in (
+            ("max_active_tasks", conf.getint, "core", 
"max_active_tasks_per_dag"),
+            ("max_active_runs", conf.getint, "core", 
"max_active_runs_per_dag"),
+            (
+                "max_consecutive_failed_dag_runs",
+                conf.getint,
+                "core",
+                "max_consecutive_failed_dag_runs_per_dag",
+            ),
+            ("catchup", conf.getboolean, "scheduler", "catchup_by_default"),
+            ("disable_bundle_versioning", conf.getboolean, "dag_processor", 
"disable_bundle_versioning"),
+        ):
+            if key not in dag:
+                dag[key] = get(section, option)
+
+    @classmethod
+    def validate_serialized_dag(cls, serialized_obj: dict[str, Any]) -> None:
+        """
+        Check that a serialized Dag, such as one a Lang-SDK runtime produced, 
can be stored and loaded.
+
+        It must match the JSON schema, deserialize, and have no cycle in its 
task graph.
+        *serialized_obj* is not changed.
+
+        :raises DeserializationError: if it does not.
+        """
+        from jsonschema import ValidationError
+
+        dag = serialized_obj.get("dag")
+        dag_id = dag.get("dag_id") if isinstance(dag, dict) else None
+        try:
+            cls.validate_schema(serialized_obj)
+        except ValidationError as e:
+            raise DeserializationError(
+                dag_id, f"Dag {dag_id!r} does not match the schema: 
{e.message}"

Review Comment:
   `e.message` is only the leaf message, so a runtime author sees `'many' is 
not of type 'number'` with no field name. Adding `e.json_path` (here 
`$.dag.max_active_runs`) would point at the field.



##########
airflow-core/src/airflow/serialization/serialized_objects.py:
##########
@@ -2131,6 +2134,70 @@ def from_dict(cls, serialized_obj: dict) -> 
SerializedDAG:
         # Pass client_defaults directly to deserialize_dag
         return cls.deserialize_dag(serialized_obj["dag"], client_defaults)
 
+    @classmethod
+    def fill_config_defaults(cls, serialized_obj: dict[str, Any]) -> None:
+        """
+        Fill in the Dag settings a serialized Dag leaves unset from the 
Airflow config, as a Python Dag does.
+
+        A Lang-SDK runtime cannot read the Airflow config, so it leaves 
``max_active_tasks``,
+        ``max_active_runs``, ``max_consecutive_failed_dag_runs``, ``catchup`` 
and
+        ``disable_bundle_versioning`` out unless the Dag sets them. A value 
the Dag sets is kept.
+        *serialized_obj* is changed in place.
+        """
+        dag = serialized_obj.get("dag")
+        if not isinstance(dag, dict):
+            # validate_serialized_dag rejects it.
+            return
+        for key, get, section, option in (
+            ("max_active_tasks", conf.getint, "core", 
"max_active_tasks_per_dag"),
+            ("max_active_runs", conf.getint, "core", 
"max_active_runs_per_dag"),
+            (
+                "max_consecutive_failed_dag_runs",
+                conf.getint,
+                "core",
+                "max_consecutive_failed_dag_runs_per_dag",
+            ),
+            ("catchup", conf.getboolean, "scheduler", "catchup_by_default"),
+            ("disable_bundle_versioning", conf.getboolean, "dag_processor", 
"disable_bundle_versioning"),
+        ):
+            if key not in dag:
+                dag[key] = get(section, option)
+
+    @classmethod
+    def validate_serialized_dag(cls, serialized_obj: dict[str, Any]) -> None:
+        """
+        Check that a serialized Dag, such as one a Lang-SDK runtime produced, 
can be stored and loaded.
+
+        It must match the JSON schema, deserialize, and have no cycle in its 
task graph.
+        *serialized_obj* is not changed.
+
+        :raises DeserializationError: if it does not.
+        """
+        from jsonschema import ValidationError

Review Comment:
   Could this import go to the top of the module? `import airflow` already 
loads jsonschema (config init goes through the providers manager's schema 
validator), so keeping it local doesn't save anything.



##########
airflow-core/tests/unit/serialization/test_dag_serialization.py:
##########
@@ -5155,3 +5155,110 @@ def get_weight(self, ti):
         op = BaseOperator(task_id="empty_task", 
weight_rule=NotRegisteredPriorityWeightStrategy())
         with pytest.raises(ValueError, match="Unknown priority strategy"):
             OperatorSerialization.serialize(op)
+
+
+class TestValidateSerializedDag:
+    @staticmethod
+    def _serialize() -> dict:
+        with DAG(dag_id="checked_dag", schedule=None) as dag:
+            BaseOperator(task_id="extract") >> BaseOperator(task_id="load")
+        return DagSerialization.to_dict(dag)
+
+    def test_accepts_a_dag_that_loads(self):
+        data = self._serialize()
+        data["__version"] = 2
+        before = copy.deepcopy(data)
+
+        DagSerialization.validate_serialized_dag(data)
+
+        assert data == before
+
+    @pytest.mark.parametrize(
+        ("change", "error"),
+        [
+            pytest.param(
+                {"max_active_runs": "many"},
+                "Dag 'checked_dag' does not match the schema: 'many' is not of 
type 'number'",
+                id="schema",
+            ),
+            pytest.param(
+                {"timetable": {"__type": "no.such.Timetable", "__var": {}}},
+                "Dag 'checked_dag' cannot be deserialized: 
TimetableNotRegistered: ",
+                id="deserialize",
+            ),
+        ],
+    )
+    def test_rejects_a_dag_that_does_not_load(self, change, error):
+        data = self._serialize()
+        data["dag"].update(change)
+
+        with pytest.raises(DeserializationError, match=f"^{re.escape(error)}"):
+            DagSerialization.validate_serialized_dag(data)
+
+    def test_rejects_a_dag_with_a_cycle(self):
+        data = self._serialize()
+        load = next(task for task in data["dag"]["tasks"] if 
task["__var"]["task_id"] == "load")
+        load["__var"]["downstream_task_ids"] = ["extract"]
+
+        with pytest.raises(DeserializationError, match="^Dag 'checked_dag' has 
a cycle through task 'load'$"):
+            DagSerialization.validate_serialized_dag(data)
+
+
+class TestFillConfigDefaults:
+    CONFIG = {
+        ("core", "max_active_tasks_per_dag"): "7",
+        ("core", "max_active_runs_per_dag"): "3",
+        ("core", "max_consecutive_failed_dag_runs_per_dag"): "5",
+        ("scheduler", "catchup_by_default"): "True",

Review Comment:
   `catchup_by_default` and `disable_bundle_versioning` are both `True` here, 
so `fill_config_defaults` reading each other's option would still pass; setting 
one of them to `False` would pin both. Separately, the `deserialize` case 
raises `TimetableNotRegistered`, which `from_dict` lets through unwrapped, so 
nothing covers the `e.__cause__` unwrap. A task with 
`downstream_task_ids=["ghost"]` would: it gives `KeyError: 'ghost'` with the 
unwrap and the generic "unexpected error" message without it.



##########
airflow-core/adr/lang-sdk/0004-dag-parsing.md:
##########
@@ -216,9 +216,11 @@ The language runtime must produce a `DagFileParsingResult` 
that matches Python A
 | `start_date` | float (epoch) | if set | Unwrapped from `__type`/`__var` |
 | `end_date` | float (epoch) | if set | Unwrapped from `__type`/`__var` |
 | `tags` | list | if non-empty | Unwrapped from `__type`/`__var` |
-| `catchup` | bool | if `true` | |
-| `max_active_tasks` | int | if non-default | |
-| `max_active_runs` | int | if non-default | |
+| `catchup` | bool | if set | Airflow fills an unset field from its config |

Review Comment:
   With these rows a runtime leaves fields out, so the section's opening line 
("matches Python Airflow's DagSerialization format exactly") and the 
"byte-identical" line in the validation section no longer hold. Worth a 
sentence there saying the runtime may omit the fields marked here and Airflow 
fills them in before it loads the Dag.



-- 
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