ephraimbuddy commented on code in PR #69864: URL: https://github.com/apache/airflow/pull/69864#discussion_r4018354550
########## airflow-core/src/airflow/serialization/dag_version_diff.py: ########## @@ -0,0 +1,898 @@ +# 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( Review Comment: Fixed. All eight fields are now public: `python_callable_name`, `expand_input`, `op_kwargs_expand_input`, `resources`, `run_as_user`, `email`, `_operator_name` and `has_retry_policy`. A renamed `@task` callable now reports `/dag/tasks/*/python_callable_name` instead of an unnamed `custom_fields` record. The test is pinned to a real `to_dict()` payload as you suggested, asserting the allowlist both ways so it cannot drift again — it caught `is_stub` sitting in the same position shortly after. It is still only as complete as the operator shapes the sweep builds, so a new shape needs a new task there, not just a new entry. --- Drafted-by: Claude Code (Opus 5); reviewed by @ephraimbuddy before posting ########## airflow-core/src/airflow/serialization/dag_version_diff.py: ########## @@ -0,0 +1,898 @@ +# 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"}) +# _get_category classifies task fields with these; every name must be a public task field or the +# entry is unreachable, since a non-public field is aggregated under custom_fields before lookup. +_DIFF_V1_TASK_ASSET_FIELDS = frozenset({"inlets", "outlets"}) +_DIFF_V1_TASK_PARAM_FIELDS = frozenset({"params"}) +_DIFF_V1_TASK_DEPENDENCY_FIELDS = frozenset({"downstream_task_ids"}) +_DIFF_V1_TASK_METADATA_FIELDS = frozenset( + { + "doc", + "doc_json", + "doc_md", + "doc_rst", + "doc_yaml", + "owner", + "ui_color", + "ui_fgcolor", + "_task_display_name", + "task_display_name", + } +) +_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", +} +_RECURSIVE_MAPPING_PATHS = { + (), + ("dag",), + ("provenance",), + *_KEYED_COLLECTION_PATHS, +} +_DIFF_V1_TASK_GROUP_METADATA_FIELDS = frozenset( + {"group_display_name", "tooltip", "doc_md", "ui_color", "ui_fgcolor"} +) +_DIFF_V1_PUBLIC_TASK_GROUP_FIELDS = _DIFF_V1_TASK_GROUP_METADATA_FIELDS | { + "_group_id", + "prefix_group_id", + "children", + "upstream_group_ids", + "downstream_group_ids", + "upstream_task_ids", + "downstream_task_ids", + "expand_input", + "is_mapped", +} + + +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: Review Comment: Fixed at both sites: `RecursionError` is in the canonicalization tuple, and `_collect_changes` has its own handler returning `serialized_dag_recursion_limit_exceeded`. Your trace also exposed an adjacent hole. Reading `SerializedDagModel.data` decompresses and parses the row, and that read sat outside any boundary, so unparseable compressed JSON, truncated bytes (`zlib.error`) and deep nesting all escaped while the uncompressed equivalent degraded. It is now inside the boundary with a distinct `serialized_dag_decode_failed`. --- Drafted-by: Claude Code (Opus 5); reviewed by @ephraimbuddy before posting -- 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]
