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


##########
airflow-core/src/airflow/serialization/dag_version_diff.py:
##########
@@ -0,0 +1,752 @@
+# 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.
+
+"""Observed-state diffs for serialized Dag payloads."""
+
+from __future__ import annotations
+
+import copy
+import hashlib
+import json
+from collections.abc import Callable, Mapping
+from datetime import timedelta
+from enum import Enum
+from typing import Any, Literal
+
+import structlog
+
+from airflow.serialization.definitions.baseoperator import 
SerializedBaseOperator
+from airflow.serialization.definitions.mappedoperator import 
SerializedMappedOperator
+from airflow.serialization.serialized_objects import (
+    _DAG_CALLBACK_FIELDS,
+    _OPERATOR_TIMEDELTA_FIELDS,
+    DagSerialization,
+    OperatorSerialization,
+)
+
+log = structlog.get_logger(__name__)
+
+DIFF_SCHEMA_VERSION = 1
+DEFAULT_MAX_CHANGES = 500
+MAX_ALLOWED_CHANGES = 5000
+SUPPORTED_SERIALIZED_DAG_SCHEMA_VERSIONS = frozenset((1, 2, 3))

Review Comment:
   This is the only schema-derived constant in the module without a parity 
test, while `_DIFF_V1_PUBLIC_TASK_FIELDS`, `_DIFF_V1_DAG_FIELD_CATEGORIES`, 
`_DAG_CALLBACK_FIELDS` and `_OPERATOR_TIMEDELTA_FIELDS` each have one. The day 
`SERIALIZER_VERSION` becomes 4, every stored row comes back 
`unsupported_serialized_dag_schema_version:4` for every Dag in the deployment 
and the suite stays green. `frozenset(range(1, 
DagSerialization.SERIALIZER_VERSION + 1))`, or a sixth parity test alongside 
the others, would close it.



##########
airflow-core/src/airflow/serialization/dag_version_diff.py:
##########
@@ -0,0 +1,752 @@
+# 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.
+
+"""Observed-state diffs for serialized Dag payloads."""
+
+from __future__ import annotations
+
+import copy
+import hashlib
+import json
+from collections.abc import Callable, Mapping
+from datetime import timedelta
+from enum import Enum
+from typing import Any, Literal
+
+import structlog
+
+from airflow.serialization.definitions.baseoperator import 
SerializedBaseOperator
+from airflow.serialization.definitions.mappedoperator import 
SerializedMappedOperator
+from airflow.serialization.serialized_objects import (
+    _DAG_CALLBACK_FIELDS,
+    _OPERATOR_TIMEDELTA_FIELDS,
+    DagSerialization,
+    OperatorSerialization,
+)
+
+log = structlog.get_logger(__name__)
+
+DIFF_SCHEMA_VERSION = 1
+DEFAULT_MAX_CHANGES = 500
+MAX_ALLOWED_CHANGES = 5000
+SUPPORTED_SERIALIZED_DAG_SCHEMA_VERSIONS = frozenset((1, 2, 3))
+
+_ORDER_INSENSITIVE_LIST_PATHS = {
+    ("dag", "tags"),
+    ("dag", "allowed_run_types"),
+}
+_KEYED_COLLECTION_PATHS = {
+    ("dag", "tasks"),
+    ("dag", "dag_dependencies"),
+    *_ORDER_INSENSITIVE_LIST_PATHS,
+}
+_CUSTOM_TASK_FIELDS_PATH_COMPONENT = "custom_fields"
+# This allowlist is part of diff schema v1. Serializer schema changes must not
+# silently change the paths visible to callers of the diff API.
+_DIFF_V1_PUBLIC_TASK_FIELDS = frozenset(
+    {
+        "__type",
+        "_disallow_kwargs_override",
+        "_expand_input_attr",
+        "_is_mapped",
+        "_is_sensor",
+        "_logger_name",
+        "_needs_expansion",
+        "_operator_extra_links",
+        "_task_display_name",
+        "_task_module",
+        "allow_nested_operators",
+        "depends_on_past",
+        "do_xcom_push",
+        "doc",
+        "doc_json",
+        "doc_md",
+        "doc_rst",
+        "doc_yaml",
+        "downstream_task_ids",
+        "email_on_failure",
+        "email_on_retry",
+        "end_date",
+        "execution_timeout",
+        "executor",
+        "executor_config",
+        "has_on_execute_callback",
+        "has_on_failure_callback",
+        "has_on_retry_callback",
+        "has_on_skipped_callback",
+        "has_on_success_callback",
+        "ignore_first_depends_on_past",
+        "inlets",
+        "is_setup",
+        "is_teardown",
+        "map_index_template",
+        "max_active_tis_per_dag",
+        "max_active_tis_per_dagrun",
+        "max_retry_delay",
+        "multiple_outputs",
+        "on_failure_fail_dagrun",
+        "outlets",
+        "owner",
+        "params",
+        "partial_kwargs",
+        "pool",
+        "pool_slots",
+        "priority_weight",
+        "queue",
+        "render_template_as_native_obj",
+        "retries",
+        "retry_delay",
+        "retry_exponential_backoff",
+        "start_date",
+        "start_from_trigger",
+        "start_trigger_args",
+        "task_id",
+        "task_type",
+        "template_ext",
+        "template_fields",
+        "template_fields_renderers",
+        "trigger_rule",
+        "ui_color",
+        "ui_fgcolor",
+        "wait_for_downstream",
+        "wait_for_past_depends_before_skipping",
+        "weight_rule",
+    }
+)
+_DIFF_V1_REDACTED_SCHEMA_TASK_FIELDS = frozenset({"_arg_bindings"})
+_DIFF_V1_PUBLIC_PARTIAL_TASK_FIELDS = _DIFF_V1_PUBLIC_TASK_FIELDS | 
{"task_display_name"}
+# Classify every Dag schema field explicitly so new fields require a policy 
decision.
+_DIFF_V1_DAG_FIELD_CATEGORIES = {
+    "_concurrency": "schedule",
+    "_processor_dags_folder": "provenance",
+    "access_control": "authorization",
+    "allowed_run_types": "schedule",
+    "bundle_name": "provenance",
+    "catchup": "schedule",
+    "dag_dependencies": "dependency",
+    "dag_display_name": "metadata",
+    "dag_id": "metadata",
+    "dagrun_timeout": "schedule",
+    "deadline": "deadline",
+    "default_args": "param",
+    "description": "metadata",
+    "disable_bundle_versioning": "task",
+    "doc_md": "metadata",
+    "edge_info": "metadata",
+    "end_date": "schedule",
+    "fail_fast": "schedule",
+    "fileloc": "provenance",
+    "has_on_failure_callback": "callback",
+    "has_on_success_callback": "callback",
+    "is_paused_upon_creation": "schedule",
+    "max_active_runs": "schedule",
+    "max_active_tasks": "schedule",
+    "max_consecutive_failed_dag_runs": "schedule",
+    "owner_links": "metadata",
+    "params": "param",
+    "relative_fileloc": "provenance",
+    "render_template_as_native_obj": "task",
+    "rerun_with_latest_version": "task",
+    "start_date": "schedule",
+    "tags": "metadata",
+    "task_group": "task",
+    "tasks": "task",
+    "timetable": "schedule",
+    "timezone": "schedule",
+}
+_DIFF_V1_LEGACY_DAG_FIELD_CATEGORIES = {
+    "fail_stop": "schedule",
+    "on_failure_callback": "callback",
+    "on_success_callback": "callback",
+    "schedule": "schedule",
+    "schedule_interval": "schedule",
+}
+_REDACTED_RECURSIVE_MAPPING_PATHS = {
+    (),
+    ("dag",),
+    ("dag", "task_group"),
+    ("provenance",),
+    *_KEYED_COLLECTION_PATHS,
+}
+
+
+def build_unavailable_dag_diff(
+    *,
+    base_data: dict[str, Any] | None,
+    target_data: dict[str, Any] | None,
+    reason: str,
+) -> dict[str, Any]:
+    """Report a known unavailable reason without comparing the stored 
payloads."""
+    return _mark_unavailable(
+        _build_diff_result(_get_schema_version(base_data), 
_get_schema_version(target_data)), reason
+    )
+
+
+def build_serialized_dag_diff(
+    *,
+    base_data: dict[str, Any] | None,
+    target_data: dict[str, Any] | None,
+    base_provenance: Mapping[str, Any] | None = None,
+    target_provenance: Mapping[str, Any] | None = None,
+    include_values: bool = False,
+    max_changes: int = DEFAULT_MAX_CHANGES,
+) -> dict[str, Any]:
+    """
+    Build a bounded, deterministic diff from two stored serialized Dag 
payloads.
+
+    Raw values, digests, and value-derived path components are returned only 
when
+    ``include_values`` is true. Callers must authorize disclosure of the entire
+    serialized payload, including access-control role names and permission 
mappings,
+    before enabling it.
+    """
+    validate_max_changes(max_changes)
+
+    base_schema_version = _get_schema_version(base_data)
+    target_schema_version = _get_schema_version(target_data)
+    result = _build_diff_result(base_schema_version, target_schema_version)
+
+    if base_data is None or target_data is None:
+        return _mark_unavailable(result, "serialized_dag_missing")
+
+    if base_schema_version is None or target_schema_version is None:
+        return _mark_unavailable(result, 
"serialized_dag_schema_version_missing")
+
+    unsupported_versions = [
+        version
+        for version in (base_schema_version, target_schema_version)
+        if version not in SUPPORTED_SERIALIZED_DAG_SCHEMA_VERSIONS
+    ]
+    if unsupported_versions:
+        return _mark_unavailable(
+            result, 
f"unsupported_serialized_dag_schema_version:{unsupported_versions[0]}"
+        )
+
+    try:
+        base_document = _canonicalize_payload_v1(base_data)
+        target_document = _canonicalize_payload_v1(target_data)
+        base_document["provenance"] = _canonicalize_value(dict(base_provenance 
or {}), path=("provenance",))
+        target_document["provenance"] = _canonicalize_value(
+            dict(target_provenance or {}), path=("provenance",)
+        )
+    except (AttributeError, KeyError, OverflowError, TypeError, ValueError) as 
error:
+        log.warning(
+            "Serialized Dag diff canonicalization failed",
+            error_type=type(error).__name__,
+            base_schema_version=base_schema_version,
+            target_schema_version=target_schema_version,
+        )
+        return _mark_unavailable(result, 
"serialized_dag_canonicalization_failed")
+
+    collector = _ChangeCollector(max_changes=max_changes, 
include_values=include_values)
+    try:
+        _collect_changes(base_document, target_document, path=(), 
collector=collector)
+    except _JsonEncodingError:
+        log.warning(
+            "Serialized Dag diff JSON encoding failed",
+            base_schema_version=base_schema_version,
+            target_schema_version=target_schema_version,
+        )
+        return _mark_unavailable(result, "serialized_dag_json_encoding_failed")
+
+    result["changes"] = collector.changes
+    result["truncated"] = collector.is_truncated
+    return result
+
+
+class _ChangeCollector:
+    def __init__(self, *, max_changes: int, include_values: bool) -> None:
+        self.changes: list[dict[str, Any]] = []
+        self.count = 0
+        self.max_changes = max_changes
+        self.include_values = include_values
+
+    @property
+    def is_truncated(self) -> bool:
+        return self.count > self.max_changes
+
+    def add(
+        self,
+        *,
+        path: tuple[str, ...],
+        operation: Literal["added", "removed", "changed"],
+        before: Any,
+        after: Any,
+    ) -> None:
+        self.count += 1
+        if len(self.changes) >= self.max_changes:
+            return
+
+        public_path = _get_public_path(path)
+        category = _get_category(public_path)
+        change: dict[str, Any] = {
+            "path": _format_path(path if self.include_values else public_path),
+            "operation": operation,
+            "category": category,
+            "impact": _get_impact(category),
+        }
+        if self.include_values:
+            change["before_digest"] = None if before is _MISSING else 
_get_digest(before)
+            change["after_digest"] = None if after is _MISSING else 
_get_digest(after)
+            if before is not _MISSING:
+                change["before_value"] = before
+            if after is not _MISSING:
+                change["after_value"] = after
+        self.changes.append(change)
+
+
+_MISSING = object()
+
+
+def validate_max_changes(max_changes: int) -> None:
+    if max_changes < 1:
+        raise ValueError("max_changes must be a positive integer")
+    if max_changes > MAX_ALLOWED_CHANGES:
+        raise ValueError(f"max_changes must not exceed {MAX_ALLOWED_CHANGES}")
+
+
+def _get_schema_version(data: Mapping[str, Any] | None) -> int | None:
+    if not isinstance(data, Mapping):
+        return None
+    version = data.get("__version")
+    return version if isinstance(version, int) and not isinstance(version, 
bool) else None
+
+
+def _mark_unavailable(result: dict[str, Any], reason: str) -> dict[str, Any]:
+    result["mode"] = "unavailable"
+    result["unavailable_reason"] = reason
+    return result
+
+
+def _build_diff_result(base_schema_version: int | None, target_schema_version: 
int | None) -> dict[str, Any]:
+    return {
+        "diff_schema_version": DIFF_SCHEMA_VERSION,
+        "serialized_dag_schema_versions": {
+            "base": base_schema_version,
+            "target": target_schema_version,
+        },
+        "mode": "observed_state",
+        "changes": [],
+        "truncated": False,
+    }
+
+
+def _canonicalize_payload_v1(data: dict[str, Any]) -> dict[str, Any]:
+    payload = copy.deepcopy(data)
+    version = _get_schema_version(payload)
+    if version is None:
+        raise ValueError("missing or invalid __version")
+    if version == 1:
+        DagSerialization.conversion_v1_to_v2(payload)
+        DagSerialization.conversion_v2_to_v3(payload)
+    elif version == 2:
+        DagSerialization.conversion_v2_to_v3(payload)
+    if not isinstance(payload.get("dag"), Mapping):
+        raise ValueError("missing dag object")
+    dag_defaults = {
+        field: value
+        for field, value in DagSerialization.get_schema_defaults("dag").items()
+        # Dag callback flags are enabled by their presence, even when their 
value is false.
+        if field not in _DAG_CALLBACK_FIELDS
+    }
+    payload["dag"] = {**dag_defaults, **payload["dag"]}
+    for field in _DAG_CALLBACK_FIELDS & payload["dag"].keys():
+        payload["dag"][field] = True
+    _apply_task_defaults(payload)
+    payload.pop("__version", None)
+    return _canonicalize_value(payload, path=())
+
+
+def _apply_task_defaults(payload: dict[str, Any]) -> None:
+    client_defaults = payload.pop("client_defaults", None)
+    if client_defaults is None:
+        client_defaults = {}
+    if not isinstance(client_defaults, Mapping):
+        raise ValueError("client_defaults is not an object")
+
+    task_defaults = client_defaults.get("tasks", {})
+    if not isinstance(task_defaults, Mapping):
+        raise ValueError("client_defaults.tasks is not an object")
+
+    schema_defaults = DagSerialization.get_schema_defaults("operator")
+    partial_fields = (
+        SerializedBaseOperator.get_serialized_fields() - 
SerializedMappedOperator.get_serialized_fields()
+    )
+    # A mapped task resolves these through partial_kwargs, so their defaults 
have to be applied
+    # there and encoded like any other partial value rather than left at the 
outer level.
+    partial_schema_defaults = {
+        field: value for field, value in schema_defaults.items() if field in 
partial_fields
+    }
+    outer_schema_defaults = {
+        field: value for field, value in schema_defaults.items() if field not 
in partial_fields
+    }
+    tasks = payload["dag"].get("tasks", [])
+    if not isinstance(tasks, list):
+        raise ValueError("dag.tasks is not a list")
+    for task in tasks:
+        if not isinstance(task, dict) or not isinstance(task.get("__var"), 
Mapping):
+            raise ValueError("task entry is not an object")
+        encoded_task = OperatorSerialization._apply_defaults_to_encoded_op(
+            dict(task["__var"]), dict(client_defaults)
+        )
+        upgraded_task = 
OperatorSerialization._upgrade_encoded_operator(encoded_task)
+        if upgraded_task.get("_is_mapped"):
+            task_data = {**outer_schema_defaults, **upgraded_task}
+            partial_kwargs = task_data.get("partial_kwargs", {})
+            if not isinstance(partial_kwargs, Mapping):
+                raise ValueError("partial_kwargs is not an object")
+            # populate_operator only folds client defaults into partial_kwargs 
when the payload
+            # carries the key, so an absent one leaves the top-level value as 
the effective value.
+            effective_partial_kwargs = (
+                {field: _encode_partial_field_value(field, value) for field, 
value in task_defaults.items()}
+                if "partial_kwargs" in task_data
+                else {}
+            )
+            for field, value in partial_kwargs.items():
+                if isinstance(value, Mapping) and "__type" in value and 
"__var" in value:
+                    effective_partial_kwargs[field] = value
+                else:
+                    effective_partial_kwargs[field] = 
_encode_partial_field_value(field, value)
+            # Match populate_operator: partial values take precedence over 
outer task defaults.
+            template_fields = task_data.get("template_fields", [])
+            for field in partial_fields & task_data.keys():
+                value = task_data.pop(field)
+                # Outer fields are already encoded unless template handling 
bypasses deserialization.
+                if field in template_fields:
+                    value = _encode_json_value(value)
+                effective_partial_kwargs.setdefault(field, value)
+            for field, value in partial_schema_defaults.items():
+                effective_partial_kwargs.setdefault(field, 
_encode_partial_field_value(field, value))
+            task_data["partial_kwargs"] = effective_partial_kwargs
+            _normalize_retry_backoff(effective_partial_kwargs)
+        else:
+            task_data = {**schema_defaults, **upgraded_task}
+            _normalize_retry_backoff(task_data)
+        task["__var"] = task_data
+
+
+def _encode_partial_field_value(field: str, value: Any) -> Any:
+    # Only partial kwargs and client defaults use field-specific timedelta 
deserialization.
+    if field in _OPERATOR_TIMEDELTA_FIELDS and value is not None:
+        return {"__type": "timedelta", "__var": value}
+    return _encode_json_value(value)
+
+
+def _encode_json_value(value: Any) -> Any:
+    """Encode plain JSON without invoking object serializers."""
+    if isinstance(value, Mapping):
+        return {"__type": "dict", "__var": {key: _encode_json_value(item) for 
key, item in value.items()}}
+    if isinstance(value, list):
+        return [_encode_json_value(item) for item in value]
+    return value
+
+
+def _normalize_retry_backoff(task_fields: dict[str, Any]) -> None:
+    if "retry_exponential_backoff" in task_fields:
+        value = task_fields["retry_exponential_backoff"]
+        task_fields["retry_exponential_backoff"] = 2.0 if value is True else 
float(value)
+
+
+def _canonicalize_value(value: Any, *, path: tuple[str, ...]) -> Any:
+    if isinstance(value, Mapping):
+        if path == ("dag", "deadline"):
+            if value.get("__type") == "deadline_alert":
+                value = value["__var"]
+            value = {"name": None, **value}
+            interval = value.get("interval")
+            if isinstance(interval, (int, float)) and not isinstance(interval, 
bool):
+                # Diff schema v1 uses the SDK's version-2 timedelta encoding 
for legacy seconds.
+                value["interval"] = {
+                    "__classname__": "datetime.timedelta",
+                    "__version__": 2,
+                    "__data__": timedelta(seconds=interval).total_seconds(),
+                }
+        return {
+            canonical_key: _canonicalize_value(item, path=path + 
(canonical_key,))
+            for canonical_key, item in (
+                (_canonicalize_mapping_key(key), item)
+                for key, item in sorted(value.items(), key=lambda item: 
_canonicalize_mapping_key(item[0]))
+            )
+        }
+    if isinstance(value, list):
+        canonical_values = [_canonicalize_value(item, path=path) for item in 
value]
+        if path == ("dag", "tasks"):
+            return _canonicalize_keyed_list(canonical_values, _get_task_id, 
path)
+        if path == ("dag", "dag_dependencies"):
+            return _canonicalize_keyed_list(canonical_values, 
_get_dependency_key, path)
+        if path in _ORDER_INSENSITIVE_LIST_PATHS:
+            return _canonicalize_keyed_list(canonical_values, _get_string_key, 
path)
+        return canonical_values
+    return value
+
+
+def _canonicalize_mapping_key(key: Any) -> str:
+    if isinstance(key, str) and isinstance(key, Enum):
+        return key.value
+    return str(key)
+
+
+def _canonicalize_keyed_list(
+    values: list[Any], key_getter: Callable[[Any], str], path: tuple[str, ...]
+) -> dict[str, Any]:
+    keyed_values: dict[str, Any] = {}
+    for value in values:
+        key = key_getter(value)
+        if key in keyed_values:
+            if path in _ORDER_INSENSITIVE_LIST_PATHS or (
+                path == ("dag", "dag_dependencies") and keyed_values[key] == 
value
+            ):
+                continue
+            raise ValueError(f"duplicate key {key!r} in /{'/'.join(path)}")
+        canonical_value = value
+        if path == ("dag", "tasks") and isinstance(value, Mapping) and "__var" 
in value:
+            task_value = dict(value["__var"])
+            if "__type" in value:
+                task_value["__type"] = value["__type"]
+            canonical_value = task_value
+        keyed_values[key] = canonical_value
+    return {key: keyed_values[key] for key in sorted(keyed_values)}
+
+
+def _get_task_id(task: Any) -> str:
+    if not isinstance(task, Mapping):
+        raise ValueError("task entry is not an object")
+    task_data = task.get("__var", task)
+    task_id = task_data.get("task_id") if isinstance(task_data, Mapping) else 
None
+    if not isinstance(task_id, str):
+        raise ValueError("task entry has no task_id")
+    return task_id
+
+
+def _get_string_key(value: Any) -> str:
+    if not isinstance(value, str):
+        raise ValueError("collection entry is not a string")
+    return value
+
+
+def _get_dependency_key(dependency: Any) -> str:
+    if not isinstance(dependency, Mapping):
+        raise ValueError("dependency entry is not an object")
+    components = (
+        dependency.get("dependency_type"),
+        dependency.get("dependency_id"),
+        dependency.get("source"),
+        dependency.get("target"),
+        dependency.get("label"),
+    )
+    return json.dumps(components, ensure_ascii=False, separators=(",", ":"))
+
+
+def _collect_changes(
+    before: Any,
+    after: Any,
+    *,
+    path: tuple[str, ...],
+    collector: _ChangeCollector,
+) -> None:
+    if isinstance(before, Mapping) and isinstance(after, Mapping):
+        if not collector.include_values and _is_task_mapping_path(path):
+            _collect_redacted_task_changes(before, after, path=path, 
collector=collector)
+            return
+        if not collector.include_values and not 
_should_recurse_redacted_mapping(path):

Review Comment:
   These two guards gate how deep the walk goes, not just whether values get 
attached, so `include_values` changes the change set itself rather than only 
redacting it. On two versions of a 100-task Dag where every task's 20-key 
`executor_config` changed, `max_changes=500` gives 100 changes and 
`truncated=False` when redacted, but 500 changes, `truncated=True` and only 25 
of the 100 tasks represented once values are authorized; the other 75 never 
appear at any path. So `max_changes` means two different things depending on a 
flag the client also sets, and the more-authorized caller gets the less 
complete answer, which contradicts `get_diff`'s docstring on line 271 ("leaving 
the structural diff available in redacted form"). Cheaper to settle here than 
after the endpoint in the next PR fixes it into a public API: compute the 
change set once with the redacted recursion shape, then attach values and 
digests to those same records when authorized.



##########
airflow-core/src/airflow/serialization/dag_version_diff.py:
##########
@@ -0,0 +1,752 @@
+# 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.
+
+"""Observed-state diffs for serialized Dag payloads."""
+
+from __future__ import annotations
+
+import copy
+import hashlib
+import json
+from collections.abc import Callable, Mapping
+from datetime import timedelta
+from enum import Enum
+from typing import Any, Literal
+
+import structlog
+
+from airflow.serialization.definitions.baseoperator import 
SerializedBaseOperator
+from airflow.serialization.definitions.mappedoperator import 
SerializedMappedOperator
+from airflow.serialization.serialized_objects import (
+    _DAG_CALLBACK_FIELDS,
+    _OPERATOR_TIMEDELTA_FIELDS,
+    DagSerialization,
+    OperatorSerialization,
+)
+
+log = structlog.get_logger(__name__)
+
+DIFF_SCHEMA_VERSION = 1
+DEFAULT_MAX_CHANGES = 500
+MAX_ALLOWED_CHANGES = 5000
+SUPPORTED_SERIALIZED_DAG_SCHEMA_VERSIONS = frozenset((1, 2, 3))
+
+_ORDER_INSENSITIVE_LIST_PATHS = {
+    ("dag", "tags"),
+    ("dag", "allowed_run_types"),
+}
+_KEYED_COLLECTION_PATHS = {
+    ("dag", "tasks"),
+    ("dag", "dag_dependencies"),
+    *_ORDER_INSENSITIVE_LIST_PATHS,
+}
+_CUSTOM_TASK_FIELDS_PATH_COMPONENT = "custom_fields"
+# This allowlist is part of diff schema v1. Serializer schema changes must not
+# silently change the paths visible to callers of the diff API.
+_DIFF_V1_PUBLIC_TASK_FIELDS = frozenset(
+    {
+        "__type",
+        "_disallow_kwargs_override",
+        "_expand_input_attr",
+        "_is_mapped",
+        "_is_sensor",
+        "_logger_name",
+        "_needs_expansion",
+        "_operator_extra_links",
+        "_task_display_name",
+        "_task_module",
+        "allow_nested_operators",
+        "depends_on_past",
+        "do_xcom_push",
+        "doc",
+        "doc_json",
+        "doc_md",
+        "doc_rst",
+        "doc_yaml",
+        "downstream_task_ids",
+        "email_on_failure",
+        "email_on_retry",
+        "end_date",
+        "execution_timeout",
+        "executor",
+        "executor_config",
+        "has_on_execute_callback",
+        "has_on_failure_callback",
+        "has_on_retry_callback",
+        "has_on_skipped_callback",
+        "has_on_success_callback",
+        "ignore_first_depends_on_past",
+        "inlets",
+        "is_setup",
+        "is_teardown",
+        "map_index_template",
+        "max_active_tis_per_dag",
+        "max_active_tis_per_dagrun",
+        "max_retry_delay",
+        "multiple_outputs",
+        "on_failure_fail_dagrun",
+        "outlets",
+        "owner",
+        "params",
+        "partial_kwargs",
+        "pool",
+        "pool_slots",
+        "priority_weight",
+        "queue",
+        "render_template_as_native_obj",
+        "retries",
+        "retry_delay",
+        "retry_exponential_backoff",
+        "start_date",
+        "start_from_trigger",
+        "start_trigger_args",
+        "task_id",
+        "task_type",
+        "template_ext",
+        "template_fields",
+        "template_fields_renderers",
+        "trigger_rule",
+        "ui_color",
+        "ui_fgcolor",
+        "wait_for_downstream",
+        "wait_for_past_depends_before_skipping",
+        "weight_rule",
+    }
+)
+_DIFF_V1_REDACTED_SCHEMA_TASK_FIELDS = frozenset({"_arg_bindings"})
+_DIFF_V1_PUBLIC_PARTIAL_TASK_FIELDS = _DIFF_V1_PUBLIC_TASK_FIELDS | 
{"task_display_name"}
+# Classify every Dag schema field explicitly so new fields require a policy 
decision.
+_DIFF_V1_DAG_FIELD_CATEGORIES = {
+    "_concurrency": "schedule",
+    "_processor_dags_folder": "provenance",
+    "access_control": "authorization",
+    "allowed_run_types": "schedule",
+    "bundle_name": "provenance",
+    "catchup": "schedule",
+    "dag_dependencies": "dependency",
+    "dag_display_name": "metadata",
+    "dag_id": "metadata",
+    "dagrun_timeout": "schedule",
+    "deadline": "deadline",
+    "default_args": "param",
+    "description": "metadata",
+    "disable_bundle_versioning": "task",
+    "doc_md": "metadata",
+    "edge_info": "metadata",
+    "end_date": "schedule",
+    "fail_fast": "schedule",
+    "fileloc": "provenance",
+    "has_on_failure_callback": "callback",
+    "has_on_success_callback": "callback",
+    "is_paused_upon_creation": "schedule",
+    "max_active_runs": "schedule",
+    "max_active_tasks": "schedule",
+    "max_consecutive_failed_dag_runs": "schedule",
+    "owner_links": "metadata",
+    "params": "param",
+    "relative_fileloc": "provenance",
+    "render_template_as_native_obj": "task",
+    "rerun_with_latest_version": "task",
+    "start_date": "schedule",
+    "tags": "metadata",
+    "task_group": "task",
+    "tasks": "task",
+    "timetable": "schedule",
+    "timezone": "schedule",
+}
+_DIFF_V1_LEGACY_DAG_FIELD_CATEGORIES = {
+    "fail_stop": "schedule",
+    "on_failure_callback": "callback",
+    "on_success_callback": "callback",
+    "schedule": "schedule",
+    "schedule_interval": "schedule",
+}
+_REDACTED_RECURSIVE_MAPPING_PATHS = {
+    (),
+    ("dag",),
+    ("dag", "task_group"),
+    ("provenance",),
+    *_KEYED_COLLECTION_PATHS,
+}
+
+
+def build_unavailable_dag_diff(
+    *,
+    base_data: dict[str, Any] | None,
+    target_data: dict[str, Any] | None,
+    reason: str,
+) -> dict[str, Any]:
+    """Report a known unavailable reason without comparing the stored 
payloads."""
+    return _mark_unavailable(
+        _build_diff_result(_get_schema_version(base_data), 
_get_schema_version(target_data)), reason
+    )
+
+
+def build_serialized_dag_diff(
+    *,
+    base_data: dict[str, Any] | None,
+    target_data: dict[str, Any] | None,
+    base_provenance: Mapping[str, Any] | None = None,
+    target_provenance: Mapping[str, Any] | None = None,
+    include_values: bool = False,
+    max_changes: int = DEFAULT_MAX_CHANGES,
+) -> dict[str, Any]:
+    """
+    Build a bounded, deterministic diff from two stored serialized Dag 
payloads.
+
+    Raw values, digests, and value-derived path components are returned only 
when
+    ``include_values`` is true. Callers must authorize disclosure of the entire
+    serialized payload, including access-control role names and permission 
mappings,
+    before enabling it.
+    """
+    validate_max_changes(max_changes)
+
+    base_schema_version = _get_schema_version(base_data)
+    target_schema_version = _get_schema_version(target_data)
+    result = _build_diff_result(base_schema_version, target_schema_version)
+
+    if base_data is None or target_data is None:
+        return _mark_unavailable(result, "serialized_dag_missing")
+
+    if base_schema_version is None or target_schema_version is None:
+        return _mark_unavailable(result, 
"serialized_dag_schema_version_missing")
+
+    unsupported_versions = [
+        version
+        for version in (base_schema_version, target_schema_version)
+        if version not in SUPPORTED_SERIALIZED_DAG_SCHEMA_VERSIONS
+    ]
+    if unsupported_versions:
+        return _mark_unavailable(
+            result, 
f"unsupported_serialized_dag_schema_version:{unsupported_versions[0]}"
+        )
+
+    try:
+        base_document = _canonicalize_payload_v1(base_data)
+        target_document = _canonicalize_payload_v1(target_data)
+        base_document["provenance"] = _canonicalize_value(dict(base_provenance 
or {}), path=("provenance",))
+        target_document["provenance"] = _canonicalize_value(
+            dict(target_provenance or {}), path=("provenance",)
+        )
+    except (AttributeError, KeyError, OverflowError, TypeError, ValueError) as 
error:
+        log.warning(
+            "Serialized Dag diff canonicalization failed",
+            error_type=type(error).__name__,
+            base_schema_version=base_schema_version,
+            target_schema_version=target_schema_version,
+        )
+        return _mark_unavailable(result, 
"serialized_dag_canonicalization_failed")
+
+    collector = _ChangeCollector(max_changes=max_changes, 
include_values=include_values)
+    try:
+        _collect_changes(base_document, target_document, path=(), 
collector=collector)
+    except _JsonEncodingError:
+        log.warning(
+            "Serialized Dag diff JSON encoding failed",
+            base_schema_version=base_schema_version,
+            target_schema_version=target_schema_version,
+        )
+        return _mark_unavailable(result, "serialized_dag_json_encoding_failed")
+
+    result["changes"] = collector.changes
+    result["truncated"] = collector.is_truncated
+    return result
+
+
+class _ChangeCollector:
+    def __init__(self, *, max_changes: int, include_values: bool) -> None:
+        self.changes: list[dict[str, Any]] = []
+        self.count = 0
+        self.max_changes = max_changes
+        self.include_values = include_values
+
+    @property
+    def is_truncated(self) -> bool:
+        return self.count > self.max_changes
+
+    def add(
+        self,
+        *,
+        path: tuple[str, ...],
+        operation: Literal["added", "removed", "changed"],
+        before: Any,
+        after: Any,
+    ) -> None:
+        self.count += 1
+        if len(self.changes) >= self.max_changes:
+            return
+
+        public_path = _get_public_path(path)
+        category = _get_category(public_path)
+        change: dict[str, Any] = {
+            "path": _format_path(path if self.include_values else public_path),
+            "operation": operation,
+            "category": category,
+            "impact": _get_impact(category),
+        }
+        if self.include_values:
+            change["before_digest"] = None if before is _MISSING else 
_get_digest(before)
+            change["after_digest"] = None if after is _MISSING else 
_get_digest(after)
+            if before is not _MISSING:
+                change["before_value"] = before
+            if after is not _MISSING:
+                change["after_value"] = after
+        self.changes.append(change)
+
+
+_MISSING = object()
+
+
+def validate_max_changes(max_changes: int) -> None:
+    if max_changes < 1:
+        raise ValueError("max_changes must be a positive integer")
+    if max_changes > MAX_ALLOWED_CHANGES:
+        raise ValueError(f"max_changes must not exceed {MAX_ALLOWED_CHANGES}")
+
+
+def _get_schema_version(data: Mapping[str, Any] | None) -> int | None:
+    if not isinstance(data, Mapping):
+        return None
+    version = data.get("__version")
+    return version if isinstance(version, int) and not isinstance(version, 
bool) else None
+
+
+def _mark_unavailable(result: dict[str, Any], reason: str) -> dict[str, Any]:
+    result["mode"] = "unavailable"
+    result["unavailable_reason"] = reason
+    return result
+
+
+def _build_diff_result(base_schema_version: int | None, target_schema_version: 
int | None) -> dict[str, Any]:
+    return {
+        "diff_schema_version": DIFF_SCHEMA_VERSION,
+        "serialized_dag_schema_versions": {
+            "base": base_schema_version,
+            "target": target_schema_version,
+        },
+        "mode": "observed_state",
+        "changes": [],
+        "truncated": False,

Review Comment:
   `values` is attached by `DagVersion.get_diff`, not here, so both exported 
engine functions return a result without that key. The description says every 
returned result includes `values.status` and the comment in `get_diff` calls 
that read always safe, which holds for the model helper but not for anything 
calling `build_serialized_dag_diff` or `build_unavailable_dag_diff` directly, 
which is the natural route for the endpoint in the next PR. Emitting it from 
`_build_diff_result` would keep `diff_schema_version: 1` one shape across all 
three entry points.



##########
airflow-core/src/airflow/serialization/dag_version_diff.py:
##########
@@ -0,0 +1,752 @@
+# 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.
+
+"""Observed-state diffs for serialized Dag payloads."""
+
+from __future__ import annotations
+
+import copy
+import hashlib
+import json
+from collections.abc import Callable, Mapping
+from datetime import timedelta
+from enum import Enum
+from typing import Any, Literal
+
+import structlog
+
+from airflow.serialization.definitions.baseoperator import 
SerializedBaseOperator
+from airflow.serialization.definitions.mappedoperator import 
SerializedMappedOperator
+from airflow.serialization.serialized_objects import (
+    _DAG_CALLBACK_FIELDS,
+    _OPERATOR_TIMEDELTA_FIELDS,
+    DagSerialization,
+    OperatorSerialization,
+)
+
+log = structlog.get_logger(__name__)
+
+DIFF_SCHEMA_VERSION = 1
+DEFAULT_MAX_CHANGES = 500
+MAX_ALLOWED_CHANGES = 5000
+SUPPORTED_SERIALIZED_DAG_SCHEMA_VERSIONS = frozenset((1, 2, 3))
+
+_ORDER_INSENSITIVE_LIST_PATHS = {
+    ("dag", "tags"),
+    ("dag", "allowed_run_types"),
+}
+_KEYED_COLLECTION_PATHS = {
+    ("dag", "tasks"),
+    ("dag", "dag_dependencies"),
+    *_ORDER_INSENSITIVE_LIST_PATHS,
+}
+_CUSTOM_TASK_FIELDS_PATH_COMPONENT = "custom_fields"
+# This allowlist is part of diff schema v1. Serializer schema changes must not
+# silently change the paths visible to callers of the diff API.
+_DIFF_V1_PUBLIC_TASK_FIELDS = frozenset(
+    {
+        "__type",
+        "_disallow_kwargs_override",
+        "_expand_input_attr",
+        "_is_mapped",
+        "_is_sensor",
+        "_logger_name",
+        "_needs_expansion",
+        "_operator_extra_links",
+        "_task_display_name",
+        "_task_module",
+        "allow_nested_operators",
+        "depends_on_past",
+        "do_xcom_push",
+        "doc",
+        "doc_json",
+        "doc_md",
+        "doc_rst",
+        "doc_yaml",
+        "downstream_task_ids",
+        "email_on_failure",
+        "email_on_retry",
+        "end_date",
+        "execution_timeout",
+        "executor",
+        "executor_config",
+        "has_on_execute_callback",
+        "has_on_failure_callback",
+        "has_on_retry_callback",
+        "has_on_skipped_callback",
+        "has_on_success_callback",
+        "ignore_first_depends_on_past",
+        "inlets",
+        "is_setup",
+        "is_teardown",
+        "map_index_template",
+        "max_active_tis_per_dag",
+        "max_active_tis_per_dagrun",
+        "max_retry_delay",
+        "multiple_outputs",
+        "on_failure_fail_dagrun",
+        "outlets",
+        "owner",
+        "params",
+        "partial_kwargs",
+        "pool",
+        "pool_slots",
+        "priority_weight",
+        "queue",
+        "render_template_as_native_obj",
+        "retries",
+        "retry_delay",
+        "retry_exponential_backoff",
+        "start_date",
+        "start_from_trigger",
+        "start_trigger_args",
+        "task_id",
+        "task_type",
+        "template_ext",
+        "template_fields",
+        "template_fields_renderers",
+        "trigger_rule",
+        "ui_color",
+        "ui_fgcolor",
+        "wait_for_downstream",
+        "wait_for_past_depends_before_skipping",
+        "weight_rule",
+    }
+)
+_DIFF_V1_REDACTED_SCHEMA_TASK_FIELDS = frozenset({"_arg_bindings"})
+_DIFF_V1_PUBLIC_PARTIAL_TASK_FIELDS = _DIFF_V1_PUBLIC_TASK_FIELDS | 
{"task_display_name"}
+# Classify every Dag schema field explicitly so new fields require a policy 
decision.
+_DIFF_V1_DAG_FIELD_CATEGORIES = {
+    "_concurrency": "schedule",
+    "_processor_dags_folder": "provenance",
+    "access_control": "authorization",
+    "allowed_run_types": "schedule",
+    "bundle_name": "provenance",
+    "catchup": "schedule",
+    "dag_dependencies": "dependency",
+    "dag_display_name": "metadata",
+    "dag_id": "metadata",
+    "dagrun_timeout": "schedule",
+    "deadline": "deadline",
+    "default_args": "param",
+    "description": "metadata",
+    "disable_bundle_versioning": "task",
+    "doc_md": "metadata",
+    "edge_info": "metadata",
+    "end_date": "schedule",
+    "fail_fast": "schedule",
+    "fileloc": "provenance",
+    "has_on_failure_callback": "callback",
+    "has_on_success_callback": "callback",
+    "is_paused_upon_creation": "schedule",
+    "max_active_runs": "schedule",
+    "max_active_tasks": "schedule",
+    "max_consecutive_failed_dag_runs": "schedule",
+    "owner_links": "metadata",
+    "params": "param",
+    "relative_fileloc": "provenance",
+    "render_template_as_native_obj": "task",
+    "rerun_with_latest_version": "task",
+    "start_date": "schedule",
+    "tags": "metadata",
+    "task_group": "task",
+    "tasks": "task",
+    "timetable": "schedule",
+    "timezone": "schedule",
+}
+_DIFF_V1_LEGACY_DAG_FIELD_CATEGORIES = {
+    "fail_stop": "schedule",
+    "on_failure_callback": "callback",
+    "on_success_callback": "callback",
+    "schedule": "schedule",
+    "schedule_interval": "schedule",
+}
+_REDACTED_RECURSIVE_MAPPING_PATHS = {
+    (),
+    ("dag",),
+    ("dag", "task_group"),
+    ("provenance",),
+    *_KEYED_COLLECTION_PATHS,
+}
+
+
+def build_unavailable_dag_diff(
+    *,
+    base_data: dict[str, Any] | None,
+    target_data: dict[str, Any] | None,
+    reason: str,
+) -> dict[str, Any]:
+    """Report a known unavailable reason without comparing the stored 
payloads."""
+    return _mark_unavailable(
+        _build_diff_result(_get_schema_version(base_data), 
_get_schema_version(target_data)), reason
+    )
+
+
+def build_serialized_dag_diff(
+    *,
+    base_data: dict[str, Any] | None,
+    target_data: dict[str, Any] | None,
+    base_provenance: Mapping[str, Any] | None = None,
+    target_provenance: Mapping[str, Any] | None = None,
+    include_values: bool = False,
+    max_changes: int = DEFAULT_MAX_CHANGES,
+) -> dict[str, Any]:
+    """
+    Build a bounded, deterministic diff from two stored serialized Dag 
payloads.
+
+    Raw values, digests, and value-derived path components are returned only 
when
+    ``include_values`` is true. Callers must authorize disclosure of the entire
+    serialized payload, including access-control role names and permission 
mappings,
+    before enabling it.
+    """
+    validate_max_changes(max_changes)
+
+    base_schema_version = _get_schema_version(base_data)
+    target_schema_version = _get_schema_version(target_data)
+    result = _build_diff_result(base_schema_version, target_schema_version)
+
+    if base_data is None or target_data is None:
+        return _mark_unavailable(result, "serialized_dag_missing")
+
+    if base_schema_version is None or target_schema_version is None:
+        return _mark_unavailable(result, 
"serialized_dag_schema_version_missing")
+
+    unsupported_versions = [
+        version
+        for version in (base_schema_version, target_schema_version)
+        if version not in SUPPORTED_SERIALIZED_DAG_SCHEMA_VERSIONS
+    ]
+    if unsupported_versions:
+        return _mark_unavailable(
+            result, 
f"unsupported_serialized_dag_schema_version:{unsupported_versions[0]}"
+        )
+
+    try:
+        base_document = _canonicalize_payload_v1(base_data)
+        target_document = _canonicalize_payload_v1(target_data)
+        base_document["provenance"] = _canonicalize_value(dict(base_provenance 
or {}), path=("provenance",))
+        target_document["provenance"] = _canonicalize_value(
+            dict(target_provenance or {}), path=("provenance",)
+        )
+    except (AttributeError, KeyError, OverflowError, TypeError, ValueError) as 
error:
+        log.warning(
+            "Serialized Dag diff canonicalization failed",
+            error_type=type(error).__name__,
+            base_schema_version=base_schema_version,
+            target_schema_version=target_schema_version,
+        )
+        return _mark_unavailable(result, 
"serialized_dag_canonicalization_failed")
+
+    collector = _ChangeCollector(max_changes=max_changes, 
include_values=include_values)
+    try:
+        _collect_changes(base_document, target_document, path=(), 
collector=collector)
+    except _JsonEncodingError:
+        log.warning(
+            "Serialized Dag diff JSON encoding failed",
+            base_schema_version=base_schema_version,
+            target_schema_version=target_schema_version,
+        )
+        return _mark_unavailable(result, "serialized_dag_json_encoding_failed")
+
+    result["changes"] = collector.changes
+    result["truncated"] = collector.is_truncated
+    return result
+
+
+class _ChangeCollector:
+    def __init__(self, *, max_changes: int, include_values: bool) -> None:
+        self.changes: list[dict[str, Any]] = []
+        self.count = 0
+        self.max_changes = max_changes
+        self.include_values = include_values
+
+    @property
+    def is_truncated(self) -> bool:
+        return self.count > self.max_changes
+
+    def add(
+        self,
+        *,
+        path: tuple[str, ...],
+        operation: Literal["added", "removed", "changed"],
+        before: Any,
+        after: Any,
+    ) -> None:
+        self.count += 1
+        if len(self.changes) >= self.max_changes:
+            return
+
+        public_path = _get_public_path(path)
+        category = _get_category(public_path)
+        change: dict[str, Any] = {
+            "path": _format_path(path if self.include_values else public_path),
+            "operation": operation,
+            "category": category,
+            "impact": _get_impact(category),
+        }
+        if self.include_values:
+            change["before_digest"] = None if before is _MISSING else 
_get_digest(before)
+            change["after_digest"] = None if after is _MISSING else 
_get_digest(after)
+            if before is not _MISSING:
+                change["before_value"] = before
+            if after is not _MISSING:
+                change["after_value"] = after
+        self.changes.append(change)
+
+
+_MISSING = object()
+
+
+def validate_max_changes(max_changes: int) -> None:
+    if max_changes < 1:
+        raise ValueError("max_changes must be a positive integer")
+    if max_changes > MAX_ALLOWED_CHANGES:
+        raise ValueError(f"max_changes must not exceed {MAX_ALLOWED_CHANGES}")
+
+
+def _get_schema_version(data: Mapping[str, Any] | None) -> int | None:
+    if not isinstance(data, Mapping):
+        return None
+    version = data.get("__version")
+    return version if isinstance(version, int) and not isinstance(version, 
bool) else None
+
+
+def _mark_unavailable(result: dict[str, Any], reason: str) -> dict[str, Any]:
+    result["mode"] = "unavailable"
+    result["unavailable_reason"] = reason
+    return result
+
+
+def _build_diff_result(base_schema_version: int | None, target_schema_version: 
int | None) -> dict[str, Any]:
+    return {
+        "diff_schema_version": DIFF_SCHEMA_VERSION,
+        "serialized_dag_schema_versions": {
+            "base": base_schema_version,
+            "target": target_schema_version,
+        },
+        "mode": "observed_state",
+        "changes": [],
+        "truncated": False,
+    }
+
+
+def _canonicalize_payload_v1(data: dict[str, Any]) -> dict[str, Any]:
+    payload = copy.deepcopy(data)
+    version = _get_schema_version(payload)
+    if version is None:
+        raise ValueError("missing or invalid __version")
+    if version == 1:
+        DagSerialization.conversion_v1_to_v2(payload)
+        DagSerialization.conversion_v2_to_v3(payload)
+    elif version == 2:
+        DagSerialization.conversion_v2_to_v3(payload)
+    if not isinstance(payload.get("dag"), Mapping):
+        raise ValueError("missing dag object")
+    dag_defaults = {
+        field: value
+        for field, value in DagSerialization.get_schema_defaults("dag").items()
+        # Dag callback flags are enabled by their presence, even when their 
value is false.
+        if field not in _DAG_CALLBACK_FIELDS
+    }
+    payload["dag"] = {**dag_defaults, **payload["dag"]}
+    for field in _DAG_CALLBACK_FIELDS & payload["dag"].keys():
+        payload["dag"][field] = True
+    _apply_task_defaults(payload)
+    payload.pop("__version", None)
+    return _canonicalize_value(payload, path=())
+
+
+def _apply_task_defaults(payload: dict[str, Any]) -> None:
+    client_defaults = payload.pop("client_defaults", None)
+    if client_defaults is None:
+        client_defaults = {}
+    if not isinstance(client_defaults, Mapping):
+        raise ValueError("client_defaults is not an object")
+
+    task_defaults = client_defaults.get("tasks", {})

Review Comment:
   The whole `client_defaults` section gets popped but only `tasks` is folded 
back in, so anything else in there is dropped before the comparison rather than 
rejected. Two payloads differing only in `client_defaults["dags"]["catchup"]` 
come back `observed_state` with `changes: []`. Given `client_defaults` is the 
growth mechanism v3 introduced, raising on unrecognised keys the way the two 
checks just above already raise on bad types would stop the diff going quiet 
the first time a new section lands.



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