diff --git a/airflow-core/docs/security/api_permissions_ref.rst b/airflow-core/docs/security/api_permissions_ref.rst index 11097ad5a422e..75d19ea81d2b9 100644 --- a/airflow-core/docs/security/api_permissions_ref.rst +++ b/airflow-core/docs/security/api_permissions_ref.rst @@ -458,6 +458,10 @@ source code so it stays up to date as endpoints are added or changed. - ``/api/v2/dags/{dag_id}/dagVersions`` - ``DAG.VERSION`` - ``GET`` + * - ``GET`` + - ``/api/v2/dags/{dag_id}/dagVersions/{base_version_number}/diff/{target_version_number}`` + - ``DAG.VERSION`` + - ``GET`` * - ``GET`` - ``/api/v2/dags/{dag_id}/dagVersions/{version_number}`` - ``DAG.VERSION`` diff --git a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_versions.py b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_versions.py index a9a7a585fd446..33d528e249c85 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_versions.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/datamodels/dag_versions.py @@ -18,6 +18,7 @@ from collections.abc import Iterable from datetime import datetime +from typing import Any, Literal from uuid import UUID from pydantic import AliasPath, Field @@ -44,3 +45,63 @@ class DAGVersionCollectionResponse(BaseModel): dag_versions: Iterable[DagVersionResponse] total_entries: int + + +class DagVersionDiffChange(BaseModel): + """One observed-state change between two serialized Dag versions.""" + + path: str + operation: Literal["added", "removed", "changed"] + category: Literal[ + "task", + "dependency", + "schedule", + "param", + "asset", + "callback", + "deadline", + "metadata", + "provenance", + "unknown", + ] + impact: Literal["execution", "metadata", "provenance", "unknown"] + before_digest: str | None = None + after_digest: str | None = None + before_value: Any | None = None + after_value: Any | None = None + + +class DagVersionDiffSourceSide(BaseModel): + """Source metadata for one side of a Dag version comparison.""" + + digest: str | None = None + content: str | None = None + + +class DagVersionDiffSource(BaseModel): + """Source comparison metadata.""" + + status: Literal["current_stored_code", "redacted", "unavailable"] + fidelity: Literal["current_stored_code", "redacted", "unavailable"] + changed: bool | None = None + base: DagVersionDiffSourceSide | None = None + target: DagVersionDiffSourceSide | None = None + + +class DagVersionDiffValues(BaseModel): + """Visibility metadata for raw serialized Dag values.""" + + status: Literal["available", "unavailable"] + + +class DagVersionDiffResponse(BaseModel): + """Observed-state diff response for two Dag versions.""" + + diff_schema_version: int + serialized_dag_schema_versions: dict[str, int | None] + mode: Literal["observed_state", "unavailable"] + changes: list[DagVersionDiffChange] + source: DagVersionDiffSource + values: DagVersionDiffValues | None = None + truncated: bool + unavailable_reason: str | None = None diff --git a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml index 005da0c6f4b5e..a4d40d9d2f8f3 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml +++ b/airflow-core/src/airflow/api_fastapi/core_api/openapi/v2-rest-api-generated.yaml @@ -10755,6 +10755,91 @@ paths: application/json: schema: $ref: '#/components/schemas/HTTPValidationError' + /api/v2/dags/{dag_id}/dagVersions/{base_version_number}/diff/{target_version_number}: + get: + tags: + - DagVersion + summary: Get Dag Version Diff + description: Compare two currently stored Dag versions. + operationId: get_dag_version_diff + security: + - OAuth2PasswordBearer: [] + - HTTPBearer: [] + parameters: + - name: dag_id + in: path + required: true + schema: + type: string + title: Dag Id + - name: base_version_number + in: path + required: true + schema: + type: integer + minimum: 1 + title: Base Version Number + - name: target_version_number + in: path + required: true + schema: + type: integer + minimum: 1 + title: Target Version Number + - name: include_values + in: query + required: false + schema: + type: boolean + default: false + title: Include Values + - name: include_source + in: query + required: false + schema: + type: boolean + default: false + title: Include Source + - name: max_changes + in: query + required: false + schema: + type: integer + maximum: 5000 + exclusiveMinimum: 0 + default: 500 + title: Max Changes + responses: + '200': + description: Successful Response + content: + application/json: + schema: + $ref: '#/components/schemas/DagVersionDiffResponse' + '401': + content: + application/json: + schema: + $ref: '#/components/schemas/HTTPExceptionResponse' + description: Unauthorized + '403': + content: + application/json: + schema: + $ref: '#/components/schemas/HTTPExceptionResponse' + description: Forbidden + '404': + content: + application/json: + schema: + $ref: '#/components/schemas/HTTPExceptionResponse' + description: Not Found + '422': + description: Validation Error + content: + application/json: + schema: + $ref: '#/components/schemas/HTTPValidationError' /api/v2/dags/{dag_id}/dagVersions: get: tags: @@ -14520,6 +14605,178 @@ components: - dag_display_name title: DagTagResponse description: Dag Tag serializer for responses. + DagVersionDiffChange: + properties: + path: + type: string + title: Path + operation: + type: string + enum: + - added + - removed + - changed + title: Operation + category: + type: string + enum: + - task + - dependency + - schedule + - param + - asset + - callback + - deadline + - metadata + - provenance + - unknown + title: Category + impact: + type: string + enum: + - execution + - metadata + - provenance + - unknown + title: Impact + before_digest: + anyOf: + - type: string + - type: 'null' + title: Before Digest + after_digest: + anyOf: + - type: string + - type: 'null' + title: After Digest + before_value: + anyOf: + - {} + - type: 'null' + title: Before Value + after_value: + anyOf: + - {} + - type: 'null' + title: After Value + type: object + required: + - path + - operation + - category + - impact + title: DagVersionDiffChange + description: One observed-state change between two serialized Dag versions. + DagVersionDiffResponse: + properties: + diff_schema_version: + type: integer + title: Diff Schema Version + serialized_dag_schema_versions: + additionalProperties: + anyOf: + - type: integer + - type: 'null' + type: object + title: Serialized Dag Schema Versions + mode: + type: string + enum: + - observed_state + - unavailable + title: Mode + changes: + items: + $ref: '#/components/schemas/DagVersionDiffChange' + type: array + title: Changes + source: + $ref: '#/components/schemas/DagVersionDiffSource' + values: + anyOf: + - $ref: '#/components/schemas/DagVersionDiffValues' + - type: 'null' + truncated: + type: boolean + title: Truncated + unavailable_reason: + anyOf: + - type: string + - type: 'null' + title: Unavailable Reason + type: object + required: + - diff_schema_version + - serialized_dag_schema_versions + - mode + - changes + - source + - truncated + title: DagVersionDiffResponse + description: Observed-state diff response for two Dag versions. + DagVersionDiffSource: + properties: + status: + type: string + enum: + - current_stored_code + - redacted + - unavailable + title: Status + fidelity: + type: string + enum: + - current_stored_code + - redacted + - unavailable + title: Fidelity + changed: + anyOf: + - type: boolean + - type: 'null' + title: Changed + base: + anyOf: + - $ref: '#/components/schemas/DagVersionDiffSourceSide' + - type: 'null' + target: + anyOf: + - $ref: '#/components/schemas/DagVersionDiffSourceSide' + - type: 'null' + type: object + required: + - status + - fidelity + title: DagVersionDiffSource + description: Source comparison metadata. + DagVersionDiffSourceSide: + properties: + digest: + anyOf: + - type: string + - type: 'null' + title: Digest + content: + anyOf: + - type: string + - type: 'null' + title: Content + type: object + title: DagVersionDiffSourceSide + description: Source metadata for one side of a Dag version comparison. + DagVersionDiffValues: + properties: + status: + type: string + enum: + - available + - unavailable + title: Status + type: object + required: + - status + title: DagVersionDiffValues + description: Visibility metadata for raw serialized Dag values. DagVersionResponse: properties: id: diff --git a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_versions.py b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_versions.py index d370e113b0a04..845367797261a 100644 --- a/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_versions.py +++ b/airflow-core/src/airflow/api_fastapi/core_api/routes/public/dag_versions.py @@ -16,13 +16,14 @@ # under the License. from __future__ import annotations -from typing import Annotated +from typing import TYPE_CHECKING, Annotated, Literal -from fastapi import Depends, HTTPException, status +from fastapi import Depends, HTTPException, Path, Query, status from sqlalchemy import select from sqlalchemy.orm import joinedload -from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity +from airflow.api_fastapi.app import get_auth_manager +from airflow.api_fastapi.auth.managers.models.resource_details import DagAccessEntity, DagDetails from airflow.api_fastapi.common.dagbag import DagBagDep, get_latest_version_of_dag from airflow.api_fastapi.common.db.common import SessionDep, paginated_select from airflow.api_fastapi.common.parameters import ( @@ -35,17 +36,36 @@ from airflow.api_fastapi.common.router import AirflowRouter from airflow.api_fastapi.core_api.datamodels.dag_versions import ( DAGVersionCollectionResponse, + DagVersionDiffResponse, DagVersionResponse, ) from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc from airflow.api_fastapi.core_api.security import ( + GetUserDep, ReadableDagVersionsFilterDep, requires_access_dag, ) from airflow.models.dag_version import DagVersion +from airflow.models.dag_version_diff import ( + MAX_ALLOWED_CHANGES, + DagVersionNotFoundError, + SourceStatus, + ValuesStatus, + get_dag_version_diff as build_dag_version_diff, +) +from airflow.models.dagbundle import DagBundleModel +from airflow.models.dagcode import DagCode +from airflow.models.team import Team + +if TYPE_CHECKING: + from sqlalchemy.orm import Session + + from airflow.api_fastapi.auth.managers.models.base_user import BaseUser dag_versions_router = AirflowRouter(tags=["DagVersion"], prefix="/dags/{dag_id}/dagVersions") +RawDataAuthorizationStatus = Literal["available", "redacted", "unavailable"] + @dag_versions_router.get( "/{version_number}", @@ -77,6 +97,180 @@ def get_dag_version( return dag_version +@dag_versions_router.get( + "/{base_version_number}/diff/{target_version_number}", + response_model_exclude_unset=True, + responses=create_openapi_http_exception_doc( + [ + status.HTTP_404_NOT_FOUND, + ] + ), + dependencies=[Depends(requires_access_dag(method="GET", access_entity=DagAccessEntity.VERSION))], +) +def get_dag_version_diff( + dag_id: str, + base_version_number: Annotated[int, Path(ge=1)], + target_version_number: Annotated[int, Path(ge=1)], + session: SessionDep, + user: GetUserDep, + include_values: Annotated[bool, Query()] = False, + include_source: Annotated[bool, Query()] = False, + max_changes: Annotated[int, Query(gt=0, le=MAX_ALLOWED_CHANGES)] = 500, +) -> DagVersionDiffResponse: + """Compare two currently stored Dag versions.""" + source_status: SourceStatus | None = None + values_status: ValuesStatus | None = None + if include_values or include_source: + versions_by_number = _get_requested_versions( + dag_id, + base_version_number, + target_version_number, + session=session, + ) + raw_data_authorization_status = _get_raw_data_authorization_status( + dag_id, + base_version_number, + target_version_number, + versions_by_number=versions_by_number, + user=user, + session=session, + ) + if include_values: + values_status = "available" if raw_data_authorization_status == "available" else "unavailable" + if include_source: + source_status = _get_source_status( + base_version_number, + target_version_number, + versions_by_number=versions_by_number, + authorization_status=raw_data_authorization_status, + session=session, + user=user, + ) + + try: + result = build_dag_version_diff( + dag_id, + base_version_number, + target_version_number, + include_values=include_values, + include_source=include_source, + max_changes=max_changes, + source_status=source_status, + values_status=values_status, + session=session, + ) + except DagVersionNotFoundError as error: + raise HTTPException(status.HTTP_404_NOT_FOUND, str(error)) from error + + return DagVersionDiffResponse.model_validate(result) + + +def _get_requested_versions( + dag_id: str, + base_version_number: int, + target_version_number: int, + *, + session: Session, +) -> dict[int, DagVersion]: + versions = session.scalars( + select(DagVersion).where( + DagVersion.dag_id == dag_id, + DagVersion.version_number.in_((base_version_number, target_version_number)), + ) + ).all() + return {version.version_number: version for version in versions} + + +def _get_raw_data_authorization_status( + dag_id: str, + base_version_number: int, + target_version_number: int, + *, + versions_by_number: dict[int, DagVersion], + session: Session, + user: BaseUser, +) -> RawDataAuthorizationStatus: + """Authorize raw values and source against each stored version's bundle.""" + requested_versions = [] + for version_number in (base_version_number, target_version_number): + if version := versions_by_number.get(version_number): + requested_versions.append(version) + else: + return "unavailable" + + bundle_names = {version.bundle_name for version in requested_versions} + if None in bundle_names: + return "unavailable" + concrete_bundle_names = {bundle_name for bundle_name in bundle_names if bundle_name is not None} + team_names_by_bundle = _get_bundle_team_names(concrete_bundle_names, session=session) + if concrete_bundle_names != set(team_names_by_bundle): + return "unavailable" + + auth_manager = get_auth_manager() + for version in requested_versions: + bundle_name = version.bundle_name + if bundle_name is None: + return "unavailable" + if not auth_manager.is_authorized_dag( + method="GET", + access_entity=DagAccessEntity.CODE, + details=DagDetails(id=dag_id, team_name=team_names_by_bundle[bundle_name]), + user=user, + ): + return "redacted" + return "available" + + +def _get_bundle_team_names(bundle_names: set[str], *, session: Session) -> dict[str, str | None]: + rows = session.execute( + select(DagBundleModel.name, Team.name) + .outerjoin(DagBundleModel.teams) + .where(DagBundleModel.name.in_(bundle_names)) + ).all() + team_names: dict[str, str | None] = {} + for bundle_name, team_name in rows: + team_names[bundle_name] = team_name + return team_names + + +def _get_source_status( + base_version_number: int, + target_version_number: int, + *, + versions_by_number: dict[int, DagVersion], + authorization_status: RawDataAuthorizationStatus, + session: Session, + user: BaseUser, +) -> SourceStatus: + """Return the source visibility state after checking both versions and historical co-location.""" + if authorization_status != "available": + return "redacted" if authorization_status == "redacted" else "unavailable" + + requested_versions = [] + for version_number in (base_version_number, target_version_number): + if version := versions_by_number.get(version_number): + requested_versions.append(version) + else: + return "unavailable" + + dag_codes = session.scalars( + select(DagCode).where(DagCode.dag_version_id.in_([version.id for version in requested_versions])) + ).all() + dag_codes_by_version_id = {dag_code.dag_version_id: dag_code for dag_code in dag_codes} + if any(version.id not in dag_codes_by_version_id for version in requested_versions): + return "unavailable" + + auth_manager = get_auth_manager() + readable_dag_ids = auth_manager.get_authorized_dag_ids(user=user) + filelocs = {dag_codes_by_version_id[version.id].fileloc for version in requested_versions} + colocated_dag_ids = set( + session.scalars(select(DagCode.dag_id).where(DagCode.fileloc.in_(filelocs))).all() + ) + if not colocated_dag_ids.issubset(readable_dag_ids): + return "redacted" + return "current_stored_code" + + @dag_versions_router.get( "", responses=create_openapi_http_exception_doc( diff --git a/airflow-core/src/airflow/cli/cli_config.py b/airflow-core/src/airflow/cli/cli_config.py index 46a0ebfa391f0..b3cbc98f037e5 100644 --- a/airflow-core/src/airflow/cli/cli_config.py +++ b/airflow-core/src/airflow/cli/cli_config.py @@ -245,6 +245,37 @@ def string_lower_type(val): choices=("table", "json", "yaml", "plain"), default="table", ) +ARG_DAG_VERSION_FROM = Arg( + ("--from-version",), + help="The base Dag version number", + type=positive_int(allow_zero=False), +) +ARG_DAG_VERSION_TO = Arg( + ("--to-version",), + help="The target Dag version number", + type=positive_int(allow_zero=False), +) +ARG_PREVIOUS_TO_LATEST = Arg( + ("--previous-to-latest",), + help="Compare the previous Dag version with the latest Dag version", + action="store_true", +) +ARG_INCLUDE_VALUES = Arg( + ("--include-values",), + help="Include currently stored values in changes", + action="store_true", +) +ARG_INCLUDE_SOURCE = Arg( + ("--include-source",), + help="Include currently stored source code when available", + action="store_true", +) +ARG_MAX_CHANGES = Arg( + ("--max-changes",), + help="Maximum number of changes to return (default: 500, maximum: 5000)", + type=positive_int(allow_zero=False), + default=500, +) ARG_COLOR = Arg( ("--color",), help="Do emit colored output (default: auto)", @@ -1191,7 +1222,34 @@ class GroupCommand(NamedTuple): ), ), ) +DAG_VERSION_COMMANDS = ( + ActionCommand( + name="diff", + help="Compare two stored Dag versions", + description=( + "Compare the currently stored serialized state of two Dag versions. " + "This CLI command runs with operator-level authority and is not an API-user authorization boundary." + ), + func=lazy_load_command("airflow.cli.commands.dag_command.dag_version_diff"), + args=( + ARG_DAG_ID, + ARG_DAG_VERSION_FROM, + ARG_DAG_VERSION_TO, + ARG_PREVIOUS_TO_LATEST, + ARG_INCLUDE_VALUES, + ARG_INCLUDE_SOURCE, + ARG_MAX_CHANGES, + ARG_OUTPUT, + ARG_VERBOSE, + ), + ), +) DAGS_COMMANDS = ( + GroupCommand( + name="versions", + help="Inspect Dag versions", + subcommands=DAG_VERSION_COMMANDS, + ), ActionCommand( name="details", help="Get DAG details given a DAG id", @@ -2350,7 +2408,8 @@ def _remove_dag_id_opt(command: ActionCommand): subcommands=[ _remove_dag_id_opt(sp) for sp in DAGS_COMMANDS - if sp.name in ["backfill", "list-runs", "pause", "unpause", "test"] + if isinstance(sp, ActionCommand) + and sp.name in ["backfill", "list-runs", "pause", "unpause", "test"] ], ), GroupCommand( diff --git a/airflow-core/src/airflow/cli/commands/dag_command.py b/airflow-core/src/airflow/cli/commands/dag_command.py index d2d303b498c0a..81df0a3fb90bd 100644 --- a/airflow-core/src/airflow/cli/commands/dag_command.py +++ b/airflow-core/src/airflow/cli/commands/dag_command.py @@ -28,7 +28,7 @@ import re import subprocess import sys -from typing import TYPE_CHECKING, cast +from typing import TYPE_CHECKING, Any, cast from sqlalchemy import func, select @@ -43,6 +43,8 @@ from airflow.exceptions import AirflowConfigException, AirflowException from airflow.jobs.job import Job from airflow.models import DagModel, DagRun, TaskInstance +from airflow.models.dag_version import DagVersion +from airflow.models.dag_version_diff import DagVersionNotFoundError, get_dag_version_diff from airflow.models.errors import ParseImportError from airflow.models.serialized_dag import SerializedDagModel from airflow.models.taskinstance import clear_task_instances @@ -80,6 +82,113 @@ _RUN_CHUNK_SIZE = 500 +@cli_utils.action_cli +@providers_configuration_loaded +@provide_session +def dag_version_diff(args, *, session: Session = NEW_SESSION) -> None: + """Compare the currently stored state of two Dag versions.""" + if args.previous_to_latest: + if args.from_version is not None or args.to_version is not None: + raise SystemExit("--previous-to-latest cannot be combined with --from-version or --to-version.") + version_numbers = session.scalars( + select(DagVersion.version_number) + .where(DagVersion.dag_id == args.dag_id) + .order_by(DagVersion.version_number.desc()) + .limit(2) + ).all() + if len(version_numbers) < 2: + raise SystemExit(f"Dag {args.dag_id!r} must have at least two versions for --previous-to-latest.") + base_version_number, target_version_number = version_numbers[1], version_numbers[0] + elif args.from_version is None or args.to_version is None: + raise SystemExit("Provide both --from-version and --to-version, or use --previous-to-latest.") + else: + base_version_number, target_version_number = args.from_version, args.to_version + + try: + result = get_dag_version_diff( + args.dag_id, + base_version_number, + target_version_number, + include_values=args.include_values, + include_source=args.include_source, + max_changes=args.max_changes, + session=session, + ) + except DagVersionNotFoundError as error: + raise SystemExit(str(error)) from error + except ValueError as error: + raise SystemExit(str(error)) from error + + _print_dag_version_diff( + result, + dag_id=args.dag_id, + base_version_number=base_version_number, + target_version_number=target_version_number, + output=args.output, + ) + + +def _print_dag_version_diff( + result: dict[str, Any], + *, + dag_id: str, + base_version_number: int, + target_version_number: int, + output: str, +) -> None: + """Render a Dag version diff in the requested CLI format.""" + console = AirflowConsole() + if output == "json": + console.print_as_json(result) + return + if output == "yaml": + console.print_as_yaml(result) + return + + print(f"Dag: {dag_id}") + print(f"Versions: {base_version_number} -> {target_version_number}") + print(f"Mode: {result['mode']}") + print() + print("Changes") + include_values = result.get("values", {}).get("status") == "available" + changes = [] + for change in result["changes"]: + rendered_change = { + "path": change["path"], + "operation": change["operation"], + "category": change["category"], + "impact": change["impact"], + } + if include_values: + rendered_change["before_value"] = _format_dag_version_diff_value(change, "before_value") + rendered_change["after_value"] = _format_dag_version_diff_value(change, "after_value") + changes.append(rendered_change) + if changes: + console.print_as(data=changes, output=output) + else: + print("No changes") + if result.get("truncated"): + print("\nChanges truncated at the configured maximum.") + + source = result["source"] + print(f"\nSource: {source['status']}") + if source.get("status") == "current_stored_code": + for label, key in (("Base", "base"), ("Target", "target")): + content = source.get(key, {}).get("content") + if content is not None: + print(f"\n{label} source:\n{content}") + if result.get("unavailable_reason"): + print(f"Unavailable reason: {result['unavailable_reason']}") + + +def _format_dag_version_diff_value(change: dict[str, Any], key: str) -> Any: + if key not in change: + return "" + if change[key] is None: + return "null" + return change[key] + + @deprecated_for_airflowctl("airflowctl dags trigger") @cli_utils.action_cli @providers_configuration_loaded diff --git a/airflow-core/src/airflow/models/dag_version_diff.py b/airflow-core/src/airflow/models/dag_version_diff.py new file mode 100644 index 0000000000000..3bd9612668091 --- /dev/null +++ b/airflow-core/src/airflow/models/dag_version_diff.py @@ -0,0 +1,483 @@ +# 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 stored Dag versions.""" + +from __future__ import annotations + +import copy +import hashlib +import json +from collections.abc import Callable, Mapping +from typing import TYPE_CHECKING, Any, Literal + +from sqlalchemy import select +from sqlalchemy.orm import joinedload + +from airflow.models.dag_version import DagVersion +from airflow.serialization.serialized_objects import DagSerialization +from airflow.utils.session import NEW_SESSION, provide_session + +if TYPE_CHECKING: + from sqlalchemy.orm import Session + + +DIFF_SCHEMA_VERSION = 1 +DEFAULT_MAX_CHANGES = 500 +MAX_ALLOWED_CHANGES = 5000 +SUPPORTED_SERIALIZED_DAG_SCHEMA_VERSIONS = frozenset((1, 2, 3)) + +DiffMode = Literal["observed_state", "unavailable"] +SourceStatus = Literal["current_stored_code", "redacted", "unavailable"] +ValuesStatus = Literal["available", "unavailable"] + +_ORDER_INSENSITIVE_LIST_PATHS = { + ("dag", "tags"), + ("dag", "allowed_run_types"), +} + + +class DagVersionNotFoundError(ValueError): + """Raised when one of the requested Dag versions does not exist.""" + + +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.""" + _validate_max_changes(max_changes) + + base_schema_version = _get_schema_version(base_data) + target_schema_version = _get_schema_version(target_data) + result: dict[str, Any] = { + "diff_schema_version": DIFF_SCHEMA_VERSION, + "serialized_dag_schema_versions": { + "base": base_schema_version, + "target": target_schema_version, + }, + "mode": "observed_state", + "changes": [], + "truncated": False, + } + + if base_data is None or target_data is None: + return _unavailable(result, "serialized_dag_missing") + + if base_schema_version is None or target_schema_version is None: + return _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 _unavailable(result, f"unsupported_serialized_dag_schema_version:{unsupported_versions[0]}") + + try: + base_document = _canonicalize_payload(base_data) + target_document = _canonicalize_payload(target_data) + except (AttributeError, KeyError, OverflowError, TypeError, ValueError): + return _unavailable(result, "serialized_dag_canonicalization_failed") + + base_document["provenance"] = _canonicalize_value(dict(base_provenance or {}), path=("provenance",)) + target_document["provenance"] = _canonicalize_value(dict(target_provenance or {}), path=("provenance",)) + + collector = _ChangeCollector(max_changes=max_changes, include_values=include_values) + _collect_changes(base_document, target_document, path=(), collector=collector) + result["changes"] = collector.changes + result["truncated"] = collector.count > max_changes + return result + + +@provide_session +def get_dag_version_diff( + dag_id: str, + base_version_number: int, + target_version_number: int, + *, + include_values: bool = False, + include_source: bool = False, + max_changes: int = DEFAULT_MAX_CHANGES, + source_status: SourceStatus | None = None, + values_status: ValuesStatus | None = None, + session: Session = NEW_SESSION, +) -> dict[str, Any]: + """ + Compare two versions of a Dag using their currently stored state. + + ``source_status`` is supplied by callers that have an authorization context. The CLI has + operator-level authority and leaves it unset, while API callers calculate it before this + function returns any source content. + """ + if base_version_number < 1 or target_version_number < 1: + raise ValueError("Dag version numbers must be positive integers") + _validate_max_changes(max_changes) + + query = ( + select(DagVersion) + .where( + DagVersion.dag_id == dag_id, + DagVersion.version_number.in_((base_version_number, target_version_number)), + ) + .options(joinedload(DagVersion.serialized_dag)) + ) + if include_source and source_status not in {"redacted", "unavailable"}: + query = query.options(joinedload(DagVersion.dag_code)) + + versions = {version.version_number: version for version in session.scalars(query).all()} + missing_version = next( + ( + version_number + for version_number in (base_version_number, target_version_number) + if version_number not in versions + ), + None, + ) + if missing_version is not None: + raise DagVersionNotFoundError( + f"The DagVersion with dag_id: `{dag_id}` and version_number: `{missing_version}` was not found" + ) + + base_version = versions[base_version_number] + target_version = versions[target_version_number] + effective_include_values = include_values and values_status != "unavailable" + result = build_serialized_dag_diff( + base_data=base_version.serialized_dag.data if base_version.serialized_dag else None, + target_data=target_version.serialized_dag.data if target_version.serialized_dag else None, + base_provenance=_get_provenance(base_version), + target_provenance=_get_provenance(target_version), + include_values=effective_include_values, + max_changes=max_changes, + ) + if include_values: + values_available = effective_include_values and result["mode"] == "observed_state" + result["values"] = {"status": "available" if values_available else "unavailable"} + result["source"] = _get_source_diff( + base_version, + target_version, + include_source=include_source, + source_status=source_status or ("current_stored_code" if include_source else "unavailable"), + ) + 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 + + 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 + + category = _get_category(path) + change = { + "path": _format_path(path), + "operation": operation, + "category": category, + "impact": _get_impact(category), + "before_digest": None if before is _MISSING else _get_digest(before), + "after_digest": None if after is _MISSING else _get_digest(after), + } + if self.include_values: + 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 _unavailable(result: dict[str, Any], reason: str) -> dict[str, Any]: + result["mode"] = "unavailable" + result["unavailable_reason"] = reason + return result + + +def _canonicalize_payload(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") + _apply_client_defaults(payload) + payload.pop("__version", None) + return _canonicalize_value(payload, path=()) + + +def _apply_client_defaults(payload: dict[str, Any]) -> None: + client_defaults = payload.pop("client_defaults", None) + if client_defaults is None: + return + 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") + + 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") + task_data = dict(task_defaults) + task_data.update(task["__var"]) + if isinstance(partial_kwargs := task_data.get("partial_kwargs"), Mapping): + effective_partial_kwargs = dict(task_defaults) + effective_partial_kwargs.update(partial_kwargs) + task_data["partial_kwargs"] = effective_partial_kwargs + task["__var"] = task_data + + +def _canonicalize_value(value: Any, *, path: tuple[str, ...]) -> Any: + if isinstance(value, Mapping): + return { + str(key): _canonicalize_value(item, path=path + (str(key),)) + for key, item in sorted(value.items(), key=lambda item: str(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_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 == ("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): + keys = sorted({str(key) for key in before} | {str(key) for key in after}) + for key in keys: + before_value = before.get(key, _MISSING) + after_value = after.get(key, _MISSING) + _collect_changes(before_value, after_value, path=path + (key,), collector=collector) + return + + if isinstance(before, list) and isinstance(after, list): + if before != after: + collector.add(path=path, operation="changed", before=before, after=after) + return + + if before is _MISSING: + collector.add(path=path, operation="added", before=_MISSING, after=after) + elif after is _MISSING: + collector.add(path=path, operation="removed", before=before, after=_MISSING) + elif before != after: + collector.add(path=path, operation="changed", before=before, after=after) + + +def _format_path(path: tuple[str, ...]) -> str: + return "/" + "/".join(component.replace("~", "~0").replace("/", "~1") for component in path) + + +def _get_digest(value: Any) -> str: + encoded = json.dumps(value, ensure_ascii=False, separators=(",", ":"), sort_keys=True).encode() + return f"sha256:{hashlib.sha256(encoded).hexdigest()}" + + +def _get_category(path: tuple[str, ...]) -> str: + lowered_path = tuple(component.lower() for component in path) + if lowered_path and lowered_path[0] == "provenance": + return "provenance" + if "deadline" in lowered_path: + return "deadline" + if any("callback" in component for component in lowered_path): + return "callback" + if any(component in {"asset", "assets", "inlets", "outlets"} for component in lowered_path): + return "asset" + if any(component in {"param", "params", "default_args"} for component in lowered_path): + return "param" + if any( + component + in {"dag_dependencies", "dependencies", "edge_info", "upstream_task_ids", "downstream_task_ids"} + for component in lowered_path + ): + return "dependency" + if any( + component + in { + "schedule", + "schedule_interval", + "timetable", + "catchup", + "max_active_runs", + "max_active_tasks", + "dagrun_timeout", + "fail_fast", + "fail_stop", + "allowed_run_types", + } + for component in lowered_path + ): + return "schedule" + if len(lowered_path) >= 2 and lowered_path[:2] == ("dag", "tasks"): + return "task" + if any( + component in {"tags", "description", "doc_md", "owner", "owners", "email", "fileloc"} + for component in lowered_path + ): + return "metadata" + return "unknown" + + +def _get_impact(category: str) -> str: + if category == "provenance": + return "provenance" + if category in {"task", "dependency", "schedule", "param", "asset", "deadline"}: + return "execution" + if category == "metadata": + return "metadata" + return "unknown" + + +def _get_provenance(version: DagVersion) -> dict[str, Any]: + return { + "bundle_name": version.bundle_name, + "bundle_version": version.bundle_version, + "version_data": version.version_data, + } + + +def _get_source_diff( + base_version: DagVersion, + target_version: DagVersion, + *, + include_source: bool, + source_status: SourceStatus, +) -> dict[str, Any]: + if not include_source: + return {"status": "unavailable", "fidelity": "unavailable"} + if source_status != "current_stored_code": + return {"status": source_status, "fidelity": source_status} + + base_code = base_version.dag_code + target_code = target_version.dag_code + if base_code is None or target_code is None: + return {"status": "unavailable", "fidelity": "unavailable"} + + base_source = base_code.source_code + target_source = target_code.source_code + return { + "status": "current_stored_code", + "fidelity": "current_stored_code", + "changed": base_source != target_source, + "base": {"digest": _get_source_digest(base_source), "content": base_source}, + "target": {"digest": _get_source_digest(target_source), "content": target_source}, + } + + +def _get_source_digest(source: str) -> str: + return f"sha256:{hashlib.sha256(source.encode()).hexdigest()}" diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts index a8f08d3df78c0..3d040cec6ce4b 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/common.ts @@ -859,6 +859,17 @@ export const UseDagVersionServiceGetDagVersionKeyFn = ({ dagId, versionNumber }: dagId: string; versionNumber: number; }, queryKey?: Array) => [useDagVersionServiceGetDagVersionKey, ...(queryKey ?? [{ dagId, versionNumber }])]; +export type DagVersionServiceGetDagVersionDiffDefaultResponse = Awaited>; +export type DagVersionServiceGetDagVersionDiffQueryResult = UseQueryResult; +export const useDagVersionServiceGetDagVersionDiffKey = "DagVersionServiceGetDagVersionDiff"; +export const UseDagVersionServiceGetDagVersionDiffKeyFn = ({ baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }: { + baseVersionNumber: number; + dagId: string; + includeSource?: boolean; + includeValues?: boolean; + maxChanges?: number; + targetVersionNumber: number; +}, queryKey?: Array) => [useDagVersionServiceGetDagVersionDiffKey, ...(queryKey ?? [{ baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }])]; export type DagVersionServiceGetDagVersionsDefaultResponse = Awaited>; export type DagVersionServiceGetDagVersionsQueryResult = UseQueryResult; export const useDagVersionServiceGetDagVersionsKey = "DagVersionServiceGetDagVersions"; diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts index 34eff20314070..ea709a10c3721 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/ensureQueryData.ts @@ -1739,6 +1739,27 @@ export const ensureUseDagVersionServiceGetDagVersionData = (queryClient: QueryCl versionNumber: number; }) => queryClient.ensureQueryData({ queryKey: Common.UseDagVersionServiceGetDagVersionKeyFn({ dagId, versionNumber }), queryFn: () => DagVersionService.getDagVersion({ dagId, versionNumber }) }); /** +* Get Dag Version Diff +* Compare two currently stored Dag versions. +* @param data The data for the request. +* @param data.dagId +* @param data.baseVersionNumber +* @param data.targetVersionNumber +* @param data.includeValues +* @param data.includeSource +* @param data.maxChanges +* @returns DagVersionDiffResponse Successful Response +* @throws ApiError +*/ +export const ensureUseDagVersionServiceGetDagVersionDiffData = (queryClient: QueryClient, { baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }: { + baseVersionNumber: number; + dagId: string; + includeSource?: boolean; + includeValues?: boolean; + maxChanges?: number; + targetVersionNumber: number; +}) => queryClient.ensureQueryData({ queryKey: Common.UseDagVersionServiceGetDagVersionDiffKeyFn({ baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }), queryFn: () => DagVersionService.getDagVersionDiff({ baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }) }); +/** * Get Dag Versions * Get all Dag Versions. * diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts index 170a17f9e58ae..083eb2f1ce6fd 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/prefetch.ts @@ -1739,6 +1739,27 @@ export const prefetchUseDagVersionServiceGetDagVersion = (queryClient: QueryClie versionNumber: number; }) => queryClient.prefetchQuery({ queryKey: Common.UseDagVersionServiceGetDagVersionKeyFn({ dagId, versionNumber }), queryFn: () => DagVersionService.getDagVersion({ dagId, versionNumber }) }); /** +* Get Dag Version Diff +* Compare two currently stored Dag versions. +* @param data The data for the request. +* @param data.dagId +* @param data.baseVersionNumber +* @param data.targetVersionNumber +* @param data.includeValues +* @param data.includeSource +* @param data.maxChanges +* @returns DagVersionDiffResponse Successful Response +* @throws ApiError +*/ +export const prefetchUseDagVersionServiceGetDagVersionDiff = (queryClient: QueryClient, { baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }: { + baseVersionNumber: number; + dagId: string; + includeSource?: boolean; + includeValues?: boolean; + maxChanges?: number; + targetVersionNumber: number; +}) => queryClient.prefetchQuery({ queryKey: Common.UseDagVersionServiceGetDagVersionDiffKeyFn({ baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }), queryFn: () => DagVersionService.getDagVersionDiff({ baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }) }); +/** * Get Dag Versions * Get all Dag Versions. * diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts index 7efadc2e4f8fa..f0095c3049516 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/queries.ts @@ -1739,6 +1739,27 @@ export const useDagVersionServiceGetDagVersion = , "queryKey" | "queryFn">) => useQuery({ queryKey: Common.UseDagVersionServiceGetDagVersionKeyFn({ dagId, versionNumber }, queryKey), queryFn: () => DagVersionService.getDagVersion({ dagId, versionNumber }) as TData, ...options }); /** +* Get Dag Version Diff +* Compare two currently stored Dag versions. +* @param data The data for the request. +* @param data.dagId +* @param data.baseVersionNumber +* @param data.targetVersionNumber +* @param data.includeValues +* @param data.includeSource +* @param data.maxChanges +* @returns DagVersionDiffResponse Successful Response +* @throws ApiError +*/ +export const useDagVersionServiceGetDagVersionDiff = = unknown[]>({ baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }: { + baseVersionNumber: number; + dagId: string; + includeSource?: boolean; + includeValues?: boolean; + maxChanges?: number; + targetVersionNumber: number; +}, queryKey?: TQueryKey, options?: Omit, "queryKey" | "queryFn">) => useQuery({ queryKey: Common.UseDagVersionServiceGetDagVersionDiffKeyFn({ baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }, queryKey), queryFn: () => DagVersionService.getDagVersionDiff({ baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }) as TData, ...options }); +/** * Get Dag Versions * Get all Dag Versions. * diff --git a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts index d4e694554d0ff..7d193fac8c7d7 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/queries/suspense.ts @@ -1739,6 +1739,27 @@ export const useDagVersionServiceGetDagVersionSuspense = , "queryKey" | "queryFn">) => useSuspenseQuery({ queryKey: Common.UseDagVersionServiceGetDagVersionKeyFn({ dagId, versionNumber }, queryKey), queryFn: () => DagVersionService.getDagVersion({ dagId, versionNumber }) as TData, ...options }); /** +* Get Dag Version Diff +* Compare two currently stored Dag versions. +* @param data The data for the request. +* @param data.dagId +* @param data.baseVersionNumber +* @param data.targetVersionNumber +* @param data.includeValues +* @param data.includeSource +* @param data.maxChanges +* @returns DagVersionDiffResponse Successful Response +* @throws ApiError +*/ +export const useDagVersionServiceGetDagVersionDiffSuspense = = unknown[]>({ baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }: { + baseVersionNumber: number; + dagId: string; + includeSource?: boolean; + includeValues?: boolean; + maxChanges?: number; + targetVersionNumber: number; +}, queryKey?: TQueryKey, options?: Omit, "queryKey" | "queryFn">) => useSuspenseQuery({ queryKey: Common.UseDagVersionServiceGetDagVersionDiffKeyFn({ baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }, queryKey), queryFn: () => DagVersionService.getDagVersionDiff({ baseVersionNumber, dagId, includeSource, includeValues, maxChanges, targetVersionNumber }) as TData, ...options }); +/** * Get Dag Versions * Get all Dag Versions. * diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts index e1880fdd4ea89..2a6ac2e677d29 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/schemas.gen.ts @@ -4506,6 +4506,235 @@ export const $DagTagResponse = { description: 'Dag Tag serializer for responses.' } as const; +export const $DagVersionDiffChange = { + properties: { + path: { + type: 'string', + title: 'Path' + }, + operation: { + type: 'string', + enum: ['added', 'removed', 'changed'], + title: 'Operation' + }, + category: { + type: 'string', + enum: ['task', 'dependency', 'schedule', 'param', 'asset', 'callback', 'deadline', 'metadata', 'provenance', 'unknown'], + title: 'Category' + }, + impact: { + type: 'string', + enum: ['execution', 'metadata', 'provenance', 'unknown'], + title: 'Impact' + }, + before_digest: { + anyOf: [ + { + type: 'string' + }, + { + type: 'null' + } + ], + title: 'Before Digest' + }, + after_digest: { + anyOf: [ + { + type: 'string' + }, + { + type: 'null' + } + ], + title: 'After Digest' + }, + before_value: { + anyOf: [ + {}, + { + type: 'null' + } + ], + title: 'Before Value' + }, + after_value: { + anyOf: [ + {}, + { + type: 'null' + } + ], + title: 'After Value' + } + }, + type: 'object', + required: ['path', 'operation', 'category', 'impact'], + title: 'DagVersionDiffChange', + description: 'One observed-state change between two serialized Dag versions.' +} as const; + +export const $DagVersionDiffResponse = { + properties: { + diff_schema_version: { + type: 'integer', + title: 'Diff Schema Version' + }, + serialized_dag_schema_versions: { + additionalProperties: { + anyOf: [ + { + type: 'integer' + }, + { + type: 'null' + } + ] + }, + type: 'object', + title: 'Serialized Dag Schema Versions' + }, + mode: { + type: 'string', + enum: ['observed_state', 'unavailable'], + title: 'Mode' + }, + changes: { + items: { + '$ref': '#/components/schemas/DagVersionDiffChange' + }, + type: 'array', + title: 'Changes' + }, + source: { + '$ref': '#/components/schemas/DagVersionDiffSource' + }, + values: { + anyOf: [ + { + '$ref': '#/components/schemas/DagVersionDiffValues' + }, + { + type: 'null' + } + ] + }, + truncated: { + type: 'boolean', + title: 'Truncated' + }, + unavailable_reason: { + anyOf: [ + { + type: 'string' + }, + { + type: 'null' + } + ], + title: 'Unavailable Reason' + } + }, + type: 'object', + required: ['diff_schema_version', 'serialized_dag_schema_versions', 'mode', 'changes', 'source', 'truncated'], + title: 'DagVersionDiffResponse', + description: 'Observed-state diff response for two Dag versions.' +} as const; + +export const $DagVersionDiffSource = { + properties: { + status: { + type: 'string', + enum: ['current_stored_code', 'redacted', 'unavailable'], + title: 'Status' + }, + fidelity: { + type: 'string', + enum: ['current_stored_code', 'redacted', 'unavailable'], + title: 'Fidelity' + }, + changed: { + anyOf: [ + { + type: 'boolean' + }, + { + type: 'null' + } + ], + title: 'Changed' + }, + base: { + anyOf: [ + { + '$ref': '#/components/schemas/DagVersionDiffSourceSide' + }, + { + type: 'null' + } + ] + }, + target: { + anyOf: [ + { + '$ref': '#/components/schemas/DagVersionDiffSourceSide' + }, + { + type: 'null' + } + ] + } + }, + type: 'object', + required: ['status', 'fidelity'], + title: 'DagVersionDiffSource', + description: 'Source comparison metadata.' +} as const; + +export const $DagVersionDiffSourceSide = { + properties: { + digest: { + anyOf: [ + { + type: 'string' + }, + { + type: 'null' + } + ], + title: 'Digest' + }, + content: { + anyOf: [ + { + type: 'string' + }, + { + type: 'null' + } + ], + title: 'Content' + } + }, + type: 'object', + title: 'DagVersionDiffSourceSide', + description: 'Source metadata for one side of a Dag version comparison.' +} as const; + +export const $DagVersionDiffValues = { + properties: { + status: { + type: 'string', + enum: ['available', 'unavailable'], + title: 'Status' + } + }, + type: 'object', + required: ['status'], + title: 'DagVersionDiffValues', + description: 'Visibility metadata for raw serialized Dag values.' +} as const; + export const $DagVersionResponse = { properties: { id: { diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts index 0356a4a364197..69df2cb675672 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/services.gen.ts @@ -3,7 +3,7 @@ import type { CancelablePromise } from './core/CancelablePromise'; import { OpenAPI } from './core/OpenAPI'; import { request as __request } from './core/request'; -import type { GetAssetsData, GetAssetsResponse, GetAssetAliasesData, GetAssetAliasesResponse, GetAssetAliasData, GetAssetAliasResponse, GetAssetEventsData, GetAssetEventsResponse, CreateAssetEventData, CreateAssetEventResponse, MaterializeAssetData, MaterializeAssetResponse, GetAssetQueuedEventsData, GetAssetQueuedEventsResponse, DeleteAssetQueuedEventsData, DeleteAssetQueuedEventsResponse, GetAssetData, GetAssetResponse, GetDagAssetQueuedEventsData, GetDagAssetQueuedEventsResponse, DeleteDagAssetQueuedEventsData, DeleteDagAssetQueuedEventsResponse, GetDagAssetQueuedEventData, GetDagAssetQueuedEventResponse, DeleteDagAssetQueuedEventData, DeleteDagAssetQueuedEventResponse, NextRunAssetsData, NextRunAssetsResponse2, ListBackfillsData, ListBackfillsResponse, CreateBackfillData, CreateBackfillResponse, GetBackfillData, GetBackfillResponse, PauseBackfillData, PauseBackfillResponse, UnpauseBackfillData, UnpauseBackfillResponse, CancelBackfillData, CancelBackfillResponse, CreateBackfillDryRunData, CreateBackfillDryRunResponse, ListBackfillsUiData, ListBackfillsUiResponse, DeleteConnectionData, DeleteConnectionResponse, GetConnectionData, GetConnectionResponse, PatchConnectionData, PatchConnectionResponse, GetConnectionTestData, GetConnectionTestResponse, EnqueueConnectionTestData, EnqueueConnectionTestResponse, GetConnectionsData, GetConnectionsResponse, PostConnectionData, PostConnectionResponse, BulkConnectionsData, BulkConnectionsResponse, TestConnectionData, TestConnectionResponse, CreateDefaultConnectionsResponse, HookMetaDataResponse, GetDagRunData, GetDagRunResponse, DeleteDagRunData, DeleteDagRunResponse, PatchDagRunData, PatchDagRunResponse, BulkDagRunsData, BulkDagRunsResponse, GetDagRunsData, GetDagRunsResponse, TriggerDagRunData, TriggerDagRunResponse, GetUpstreamAssetEventsData, GetUpstreamAssetEventsResponse, ClearDagRunData, ClearDagRunResponse, WaitDagRunUntilFinishedData, WaitDagRunUntilFinishedResponse, GetListDagRunsBatchData, GetListDagRunsBatchResponse, ClearDagRunsData, ClearDagRunsResponse, ClearDagRunPartitionsData, ClearDagRunPartitionsResponse, GetDagRunStatsData, GetDagRunStatsResponse, GetDagSourceData, GetDagSourceResponse, GetDagStatsData, GetDagStatsResponse, GetConfigData, GetConfigResponse, GetConfigValueData, GetConfigValueResponse, GetConfigsResponse, ListDagWarningsData, ListDagWarningsResponse, GetDagsData, GetDagsResponse, PatchDagsData, PatchDagsResponse, GetDagData, GetDagResponse, PatchDagData, PatchDagResponse, DeleteDagData, DeleteDagResponse, GetDagDetailsData, GetDagDetailsResponse, FavoriteDagData, FavoriteDagResponse, UnfavoriteDagData, UnfavoriteDagResponse, GetDagTagsData, GetDagTagsResponse, GetDagsUiData, GetDagsUiResponse, GetLatestRunInfoData, GetLatestRunInfoResponse, GetDagRunStateCountsUiData, GetDagRunStateCountsUiResponse, GetEventLogData, GetEventLogResponse, GetEventLogsData, GetEventLogsResponse, GetExtraLinksData, GetExtraLinksResponse, GetTaskInstanceData, GetTaskInstanceResponse, PatchTaskInstanceData, PatchTaskInstanceResponse, DeleteTaskInstanceData, DeleteTaskInstanceResponse, GetMappedTaskInstancesData, GetMappedTaskInstancesResponse, GetTaskInstanceDependenciesByMapIndexData, GetTaskInstanceDependenciesByMapIndexResponse, GetTaskInstanceDependenciesData, GetTaskInstanceDependenciesResponse, GetTaskInstanceTriesData, GetTaskInstanceTriesResponse, GetMappedTaskInstanceTriesData, GetMappedTaskInstanceTriesResponse, GetMappedTaskInstanceData, GetMappedTaskInstanceResponse, PatchTaskInstanceByMapIndexData, PatchTaskInstanceByMapIndexResponse, GetTaskInstancesData, GetTaskInstancesResponse, BulkTaskInstancesData, BulkTaskInstancesResponse, GetTaskInstancesBatchData, GetTaskInstancesBatchResponse, GetTaskInstanceTryDetailsData, GetTaskInstanceTryDetailsResponse, GetMappedTaskInstanceTryDetailsData, GetMappedTaskInstanceTryDetailsResponse, PostClearTaskInstancesData, PostClearTaskInstancesResponse, PatchTaskGroupInstancesData, PatchTaskGroupInstancesResponse, PatchTaskGroupInstancesDryRunData, PatchTaskGroupInstancesDryRunResponse, PatchTaskInstanceDryRunByMapIndexData, PatchTaskInstanceDryRunByMapIndexResponse, PatchTaskInstanceDryRunData, PatchTaskInstanceDryRunResponse, GetLogData, GetLogResponse, GetExternalLogUrlData, GetExternalLogUrlResponse, UpdateHitlDetailData, UpdateHitlDetailResponse, GetHitlDetailData, GetHitlDetailResponse, GetHitlDetailTryDetailData, GetHitlDetailTryDetailResponse, GetHitlDetailsData, GetHitlDetailsResponse, GetImportErrorData, GetImportErrorResponse, GetImportErrorsData, GetImportErrorsResponse, GetJobsData, GetJobsResponse, GetPluginsData, GetPluginsResponse, ImportErrorsResponse, DeletePoolData, DeletePoolResponse, GetPoolData, GetPoolResponse, PatchPoolData, PatchPoolResponse, GetPoolsData, GetPoolsResponse, PostPoolData, PostPoolResponse, BulkPoolsData, BulkPoolsResponse, GetProvidersData, GetProvidersResponse, ListAssetStateStoreData, ListAssetStateStoreResponse, ClearAssetStateStoreData, ClearAssetStateStoreResponse, GetAssetStateStoreData, GetAssetStateStoreResponse, SetAssetStateStoreData, SetAssetStateStoreResponse, DeleteAssetStateStoreData, DeleteAssetStateStoreResponse, ListTaskStateStoreData, ListTaskStateStoreResponse, ClearTaskStateStoreData, ClearTaskStateStoreResponse, GetTaskStateStoreData, GetTaskStateStoreResponse, SetTaskStateStoreData, SetTaskStateStoreResponse, PatchTaskStateStoreData, PatchTaskStateStoreResponse, DeleteTaskStateStoreData, DeleteTaskStateStoreResponse, GetXcomEntryData, GetXcomEntryResponse, UpdateXcomEntryData, UpdateXcomEntryResponse, DeleteXcomEntryData, DeleteXcomEntryResponse, GetXcomEntriesData, GetXcomEntriesResponse, CreateXcomEntryData, CreateXcomEntryResponse, GetTasksData, GetTasksResponse, GetTaskData, GetTaskResponse, DeleteVariableData, DeleteVariableResponse, GetVariableData, GetVariableResponse, PatchVariableData, PatchVariableResponse, GetVariablesData, GetVariablesResponse, PostVariableData, PostVariableResponse, BulkVariablesData, BulkVariablesResponse, ReparseDagFileData, ReparseDagFileResponse, GetDagVersionData, GetDagVersionResponse, GetDagVersionsData, GetDagVersionsResponse, GetHealthResponse, GetVersionResponse, LoginData, LoginResponse, LogoutResponse, GetAuthMenusResponse, GetCurrentUserInfoResponse, GenerateTokenData, GenerateTokenResponse2, GetPartitionedDagRunsData, GetPartitionedDagRunsResponse, GetPendingPartitionedDagRunData, GetPendingPartitionedDagRunResponse, GetDependenciesData, GetDependenciesResponse, HistoricalMetricsData, HistoricalMetricsResponse, DagStatsResponse2, GetDeadlinesData, GetDeadlinesResponse, GetDagDeadlineAlertsData, GetDagDeadlineAlertsResponse, StructureDataData, StructureDataResponse2, GetDagStructureData, GetDagStructureResponse, GetGridRunsData, GetGridRunsResponse, GetGridTiSummariesStreamData, GetGridTiSummariesStreamResponse, GetGanttDataData, GetGanttDataResponse, GetCalendarData, GetCalendarResponse, GetCalendarDeadlinesData, GetCalendarDeadlinesResponse, ListTeamsData, ListTeamsResponse } from './types.gen'; +import type { GetAssetsData, GetAssetsResponse, GetAssetAliasesData, GetAssetAliasesResponse, GetAssetAliasData, GetAssetAliasResponse, GetAssetEventsData, GetAssetEventsResponse, CreateAssetEventData, CreateAssetEventResponse, MaterializeAssetData, MaterializeAssetResponse, GetAssetQueuedEventsData, GetAssetQueuedEventsResponse, DeleteAssetQueuedEventsData, DeleteAssetQueuedEventsResponse, GetAssetData, GetAssetResponse, GetDagAssetQueuedEventsData, GetDagAssetQueuedEventsResponse, DeleteDagAssetQueuedEventsData, DeleteDagAssetQueuedEventsResponse, GetDagAssetQueuedEventData, GetDagAssetQueuedEventResponse, DeleteDagAssetQueuedEventData, DeleteDagAssetQueuedEventResponse, NextRunAssetsData, NextRunAssetsResponse2, ListBackfillsData, ListBackfillsResponse, CreateBackfillData, CreateBackfillResponse, GetBackfillData, GetBackfillResponse, PauseBackfillData, PauseBackfillResponse, UnpauseBackfillData, UnpauseBackfillResponse, CancelBackfillData, CancelBackfillResponse, CreateBackfillDryRunData, CreateBackfillDryRunResponse, ListBackfillsUiData, ListBackfillsUiResponse, DeleteConnectionData, DeleteConnectionResponse, GetConnectionData, GetConnectionResponse, PatchConnectionData, PatchConnectionResponse, GetConnectionTestData, GetConnectionTestResponse, EnqueueConnectionTestData, EnqueueConnectionTestResponse, GetConnectionsData, GetConnectionsResponse, PostConnectionData, PostConnectionResponse, BulkConnectionsData, BulkConnectionsResponse, TestConnectionData, TestConnectionResponse, CreateDefaultConnectionsResponse, HookMetaDataResponse, GetDagRunData, GetDagRunResponse, DeleteDagRunData, DeleteDagRunResponse, PatchDagRunData, PatchDagRunResponse, BulkDagRunsData, BulkDagRunsResponse, GetDagRunsData, GetDagRunsResponse, TriggerDagRunData, TriggerDagRunResponse, GetUpstreamAssetEventsData, GetUpstreamAssetEventsResponse, ClearDagRunData, ClearDagRunResponse, WaitDagRunUntilFinishedData, WaitDagRunUntilFinishedResponse, GetListDagRunsBatchData, GetListDagRunsBatchResponse, ClearDagRunsData, ClearDagRunsResponse, ClearDagRunPartitionsData, ClearDagRunPartitionsResponse, GetDagRunStatsData, GetDagRunStatsResponse, GetDagSourceData, GetDagSourceResponse, GetDagStatsData, GetDagStatsResponse, GetConfigData, GetConfigResponse, GetConfigValueData, GetConfigValueResponse, GetConfigsResponse, ListDagWarningsData, ListDagWarningsResponse, GetDagsData, GetDagsResponse, PatchDagsData, PatchDagsResponse, GetDagData, GetDagResponse, PatchDagData, PatchDagResponse, DeleteDagData, DeleteDagResponse, GetDagDetailsData, GetDagDetailsResponse, FavoriteDagData, FavoriteDagResponse, UnfavoriteDagData, UnfavoriteDagResponse, GetDagTagsData, GetDagTagsResponse, GetDagsUiData, GetDagsUiResponse, GetLatestRunInfoData, GetLatestRunInfoResponse, GetDagRunStateCountsUiData, GetDagRunStateCountsUiResponse, GetEventLogData, GetEventLogResponse, GetEventLogsData, GetEventLogsResponse, GetExtraLinksData, GetExtraLinksResponse, GetTaskInstanceData, GetTaskInstanceResponse, PatchTaskInstanceData, PatchTaskInstanceResponse, DeleteTaskInstanceData, DeleteTaskInstanceResponse, GetMappedTaskInstancesData, GetMappedTaskInstancesResponse, GetTaskInstanceDependenciesByMapIndexData, GetTaskInstanceDependenciesByMapIndexResponse, GetTaskInstanceDependenciesData, GetTaskInstanceDependenciesResponse, GetTaskInstanceTriesData, GetTaskInstanceTriesResponse, GetMappedTaskInstanceTriesData, GetMappedTaskInstanceTriesResponse, GetMappedTaskInstanceData, GetMappedTaskInstanceResponse, PatchTaskInstanceByMapIndexData, PatchTaskInstanceByMapIndexResponse, GetTaskInstancesData, GetTaskInstancesResponse, BulkTaskInstancesData, BulkTaskInstancesResponse, GetTaskInstancesBatchData, GetTaskInstancesBatchResponse, GetTaskInstanceTryDetailsData, GetTaskInstanceTryDetailsResponse, GetMappedTaskInstanceTryDetailsData, GetMappedTaskInstanceTryDetailsResponse, PostClearTaskInstancesData, PostClearTaskInstancesResponse, PatchTaskGroupInstancesData, PatchTaskGroupInstancesResponse, PatchTaskGroupInstancesDryRunData, PatchTaskGroupInstancesDryRunResponse, PatchTaskInstanceDryRunByMapIndexData, PatchTaskInstanceDryRunByMapIndexResponse, PatchTaskInstanceDryRunData, PatchTaskInstanceDryRunResponse, GetLogData, GetLogResponse, GetExternalLogUrlData, GetExternalLogUrlResponse, UpdateHitlDetailData, UpdateHitlDetailResponse, GetHitlDetailData, GetHitlDetailResponse, GetHitlDetailTryDetailData, GetHitlDetailTryDetailResponse, GetHitlDetailsData, GetHitlDetailsResponse, GetImportErrorData, GetImportErrorResponse, GetImportErrorsData, GetImportErrorsResponse, GetJobsData, GetJobsResponse, GetPluginsData, GetPluginsResponse, ImportErrorsResponse, DeletePoolData, DeletePoolResponse, GetPoolData, GetPoolResponse, PatchPoolData, PatchPoolResponse, GetPoolsData, GetPoolsResponse, PostPoolData, PostPoolResponse, BulkPoolsData, BulkPoolsResponse, GetProvidersData, GetProvidersResponse, ListAssetStateStoreData, ListAssetStateStoreResponse, ClearAssetStateStoreData, ClearAssetStateStoreResponse, GetAssetStateStoreData, GetAssetStateStoreResponse, SetAssetStateStoreData, SetAssetStateStoreResponse, DeleteAssetStateStoreData, DeleteAssetStateStoreResponse, ListTaskStateStoreData, ListTaskStateStoreResponse, ClearTaskStateStoreData, ClearTaskStateStoreResponse, GetTaskStateStoreData, GetTaskStateStoreResponse, SetTaskStateStoreData, SetTaskStateStoreResponse, PatchTaskStateStoreData, PatchTaskStateStoreResponse, DeleteTaskStateStoreData, DeleteTaskStateStoreResponse, GetXcomEntryData, GetXcomEntryResponse, UpdateXcomEntryData, UpdateXcomEntryResponse, DeleteXcomEntryData, DeleteXcomEntryResponse, GetXcomEntriesData, GetXcomEntriesResponse, CreateXcomEntryData, CreateXcomEntryResponse, GetTasksData, GetTasksResponse, GetTaskData, GetTaskResponse, DeleteVariableData, DeleteVariableResponse, GetVariableData, GetVariableResponse, PatchVariableData, PatchVariableResponse, GetVariablesData, GetVariablesResponse, PostVariableData, PostVariableResponse, BulkVariablesData, BulkVariablesResponse, ReparseDagFileData, ReparseDagFileResponse, GetDagVersionData, GetDagVersionResponse, GetDagVersionDiffData, GetDagVersionDiffResponse, GetDagVersionsData, GetDagVersionsResponse, GetHealthResponse, GetVersionResponse, LoginData, LoginResponse, LogoutResponse, GetAuthMenusResponse, GetCurrentUserInfoResponse, GenerateTokenData, GenerateTokenResponse2, GetPartitionedDagRunsData, GetPartitionedDagRunsResponse, GetPendingPartitionedDagRunData, GetPendingPartitionedDagRunResponse, GetDependenciesData, GetDependenciesResponse, HistoricalMetricsData, HistoricalMetricsResponse, DagStatsResponse2, GetDeadlinesData, GetDeadlinesResponse, GetDagDeadlineAlertsData, GetDagDeadlineAlertsResponse, StructureDataData, StructureDataResponse2, GetDagStructureData, GetDagStructureResponse, GetGridRunsData, GetGridRunsResponse, GetGridTiSummariesStreamData, GetGridTiSummariesStreamResponse, GetGanttDataData, GetGanttDataResponse, GetCalendarData, GetCalendarResponse, GetCalendarDeadlinesData, GetCalendarDeadlinesResponse, ListTeamsData, ListTeamsResponse } from './types.gen'; export class AssetService { /** @@ -4524,6 +4524,42 @@ export class DagVersionService { }); } + /** + * Get Dag Version Diff + * Compare two currently stored Dag versions. + * @param data The data for the request. + * @param data.dagId + * @param data.baseVersionNumber + * @param data.targetVersionNumber + * @param data.includeValues + * @param data.includeSource + * @param data.maxChanges + * @returns DagVersionDiffResponse Successful Response + * @throws ApiError + */ + public static getDagVersionDiff(data: GetDagVersionDiffData): CancelablePromise { + return __request(OpenAPI, { + method: 'GET', + url: '/api/v2/dags/{dag_id}/dagVersions/{base_version_number}/diff/{target_version_number}', + path: { + dag_id: data.dagId, + base_version_number: data.baseVersionNumber, + target_version_number: data.targetVersionNumber + }, + query: { + include_values: data.includeValues, + include_source: data.includeSource, + max_changes: data.maxChanges + }, + errors: { + 401: 'Unauthorized', + 403: 'Forbidden', + 404: 'Not Found', + 422: 'Validation Error' + } + }); + } + /** * Get Dag Versions * Get all Dag Versions. diff --git a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts index 750202cd7fb0b..8d50f03f06048 100644 --- a/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts +++ b/airflow-core/src/airflow/ui/openapi-gen/requests/types.gen.ts @@ -1196,6 +1196,76 @@ export type DagTagResponse = { dag_display_name: string; }; +/** + * One observed-state change between two serialized Dag versions. + */ +export type DagVersionDiffChange = { + path: string; + operation: 'added' | 'removed' | 'changed'; + category: 'task' | 'dependency' | 'schedule' | 'param' | 'asset' | 'callback' | 'deadline' | 'metadata' | 'provenance' | 'unknown'; + impact: 'execution' | 'metadata' | 'provenance' | 'unknown'; + before_digest?: string | null; + after_digest?: string | null; + before_value?: unknown | null; + after_value?: unknown | null; +}; + +export type operation = 'added' | 'removed' | 'changed'; + +export type category = 'task' | 'dependency' | 'schedule' | 'param' | 'asset' | 'callback' | 'deadline' | 'metadata' | 'provenance' | 'unknown'; + +export type impact = 'execution' | 'metadata' | 'provenance' | 'unknown'; + +/** + * Observed-state diff response for two Dag versions. + */ +export type DagVersionDiffResponse = { + diff_schema_version: number; + serialized_dag_schema_versions: { + [key: string]: (number | null); + }; + mode: 'observed_state' | 'unavailable'; + changes: Array; + source: DagVersionDiffSource; + values?: DagVersionDiffValues | null; + truncated: boolean; + unavailable_reason?: string | null; +}; + +export type mode = 'observed_state' | 'unavailable'; + +/** + * Source comparison metadata. + */ +export type DagVersionDiffSource = { + status: 'current_stored_code' | 'redacted' | 'unavailable'; + fidelity: 'current_stored_code' | 'redacted' | 'unavailable'; + changed?: boolean | null; + base?: DagVersionDiffSourceSide | null; + target?: DagVersionDiffSourceSide | null; +}; + +export type status = 'current_stored_code' | 'redacted' | 'unavailable'; + +export type fidelity = 'current_stored_code' | 'redacted' | 'unavailable'; + +/** + * Source metadata for one side of a Dag version comparison. + */ +export type DagVersionDiffSourceSide = { + digest?: string | null; + content?: string | null; +}; + +/** + * Visibility metadata for raw serialized Dag values. + */ +export type DagVersionDiffValues = { + status: 'available' | 'unavailable'; +}; + +export type status2 = 'available' | 'unavailable'; + /** * Dag Version serializer for responses. */ @@ -2790,7 +2860,7 @@ export type UIAlert = { category: 'info' | 'warning' | 'error'; }; -export type category = 'info' | 'warning' | 'error'; +export type category2 = 'info' | 'warning' | 'error'; export type GetAssetsData = { dagIds?: Array<(string)>; @@ -4525,6 +4595,17 @@ export type GetDagVersionData = { export type GetDagVersionResponse = DagVersionResponse; +export type GetDagVersionDiffData = { + baseVersionNumber: number; + dagId: string; + includeSource?: boolean; + includeValues?: boolean; + maxChanges?: number; + targetVersionNumber: number; +}; + +export type GetDagVersionDiffResponse = DagVersionDiffResponse; + export type GetDagVersionsData = { bundleName?: string; bundleVersion?: string | null; @@ -8237,6 +8318,33 @@ export type $OpenApiTs = { }; }; }; + '/api/v2/dags/{dag_id}/dagVersions/{base_version_number}/diff/{target_version_number}': { + get: { + req: GetDagVersionDiffData; + res: { + /** + * Successful Response + */ + 200: DagVersionDiffResponse; + /** + * Unauthorized + */ + 401: HTTPExceptionResponse; + /** + * Forbidden + */ + 403: HTTPExceptionResponse; + /** + * Not Found + */ + 404: HTTPExceptionResponse; + /** + * Validation Error + */ + 422: HTTPValidationError; + }; + }; + }; '/api/v2/dags/{dag_id}/dagVersions': { get: { req: GetDagVersionsData; diff --git a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_versions.py b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_versions.py index 73051be51df5a..edc1a25f4bad8 100644 --- a/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_versions.py +++ b/airflow-core/tests/unit/api_fastapi/core_api/routes/public/test_dag_versions.py @@ -19,7 +19,11 @@ from unittest import mock import pytest +from sqlalchemy import event, select, update +from airflow import settings +from airflow.models.dag_version import DagVersion +from airflow.models.serialized_dag import SerializedDagModel from airflow.providers.standard.operators.empty import EmptyOperator from tests_common.test_utils.asserts import assert_queries_count @@ -210,6 +214,269 @@ def test_should_respond_403(self, unauthorized_test_client): assert response.status_code == 403 +class TestGetDagVersionDiff(TestDagVersionEndpoint): + @pytest.mark.usefixtures("make_dag_with_multiple_versions") + def test_get_dag_version_diff(self, test_client): + response = test_client.get( + "/dags/dag_with_multiple_versions/dagVersions/1/diff/3", + params={"include_values": True}, + ) + + assert response.status_code == 200 + response_data = response.json() + assert response_data["diff_schema_version"] == 1 + assert response_data["serialized_dag_schema_versions"] == {"base": 3, "target": 3} + assert response_data["mode"] == "observed_state" + assert response_data["truncated"] is False + assert response_data["values"] == {"status": "available"} + assert response_data["source"] == {"status": "unavailable", "fidelity": "unavailable"} + assert any( + change["path"] == "/dag/tasks/task2" + and change["operation"] == "added" + and change["category"] == "task" + for change in response_data["changes"] + ) + assert any("after_value" in change for change in response_data["changes"]) + + @pytest.mark.usefixtures("make_dag_with_multiple_versions") + @mock.patch("airflow.api_fastapi.core_api.routes.public.dag_versions.get_auth_manager", autospec=True) + def test_get_dag_version_diff_redacts_values_without_code_access( + self, mock_get_auth_manager, test_client + ): + mock_get_auth_manager.return_value.is_authorized_dag.return_value = False + + response = test_client.get( + "/dags/dag_with_multiple_versions/dagVersions/1/diff/3", + params={"include_values": True}, + ) + + assert response.status_code == 200 + response_data = response.json() + assert response_data["values"] == {"status": "unavailable"} + assert all( + "before_value" not in change and "after_value" not in change + for change in response_data["changes"] + ) + + @pytest.mark.usefixtures("make_dag_with_multiple_versions") + @mock.patch( + "airflow.api_fastapi.core_api.routes.public.dag_versions.build_dag_version_diff", autospec=True + ) + def test_get_dag_version_diff_preserves_explicit_null_values( + self, mock_build_dag_version_diff, test_client + ): + mock_build_dag_version_diff.return_value = { + "diff_schema_version": 1, + "serialized_dag_schema_versions": {"base": 3, "target": 3}, + "mode": "observed_state", + "changes": [ + { + "path": "/dag/value", + "operation": "added", + "category": "unknown", + "impact": "unknown", + "before_digest": None, + "after_digest": "sha256:null", + "after_value": None, + } + ], + "source": {"status": "unavailable", "fidelity": "unavailable"}, + "values": {"status": "available"}, + "truncated": False, + } + + response = test_client.get( + "/dags/dag_with_multiple_versions/dagVersions/1/diff/3", + params={"include_values": True}, + ) + + assert response.status_code == 200 + change = response.json()["changes"][0] + assert "before_value" not in change + assert "after_value" in change + assert change["after_value"] is None + + @pytest.mark.usefixtures("make_dag_with_multiple_versions") + def test_get_dag_version_diff_truncates(self, test_client): + response = test_client.get( + "/dags/dag_with_multiple_versions/dagVersions/1/diff/3", + params={"max_changes": 1}, + ) + + assert response.status_code == 200 + response_data = response.json() + assert response_data["truncated"] is True + assert "values" not in response_data + + @pytest.mark.usefixtures("make_dag_with_multiple_versions") + def test_get_dag_version_diff_marks_values_unavailable_when_diff_is_unavailable( + self, test_client, session + ): + serialized_dag = session.scalar( + select(SerializedDagModel) + .join(DagVersion) + .where( + DagVersion.dag_id == "dag_with_multiple_versions", + DagVersion.version_number == 1, + ) + ) + assert serialized_dag is not None + serialized_data = serialized_dag.data + assert serialized_data is not None + session.execute( + update(SerializedDagModel) + .where(SerializedDagModel.id == serialized_dag.id) + .values( + { + SerializedDagModel._data: {**serialized_data, "__version": 99}, + SerializedDagModel._data_compressed: None, + } + ) + ) + session.commit() + session.expunge_all() + + response = test_client.get( + "/dags/dag_with_multiple_versions/dagVersions/1/diff/3", + params={"include_values": True}, + ) + + assert response.status_code == 200 + response_data = response.json() + assert response_data["mode"] == "unavailable" + assert response_data["values"] == {"status": "unavailable"} + + def test_get_dag_version_diff_rejects_excessive_change_limit(self, test_client): + response = test_client.get( + "/dags/dag_with_multiple_versions/dagVersions/1/diff/3", + params={"max_changes": 5001}, + ) + + assert response.status_code == 422 + + @pytest.mark.usefixtures("make_dag_with_multiple_versions") + def test_get_dag_version_diff_includes_current_source(self, test_client): + response = test_client.get( + "/dags/dag_with_multiple_versions/dagVersions/1/diff/3", + params={"include_source": True}, + ) + + assert response.status_code == 200 + source = response.json()["source"] + assert source["status"] == "current_stored_code" + assert source["fidelity"] == "current_stored_code" + assert source["changed"] is False + assert source["base"]["content"] + assert source["base"]["content"] == source["target"]["content"] + assert source["base"]["digest"].startswith("sha256:") + + @pytest.mark.usefixtures("make_dag_with_multiple_versions") + @mock.patch("airflow.api_fastapi.core_api.routes.public.dag_versions.get_auth_manager", autospec=True) + def test_get_dag_version_diff_redacts_unreadable_colocated_source( + self, mock_get_auth_manager, test_client + ): + mock_get_auth_manager.return_value.is_authorized_dag.return_value = True + mock_get_auth_manager.return_value.get_authorized_dag_ids.return_value = set() + + response = test_client.get( + "/dags/dag_with_multiple_versions/dagVersions/1/diff/3", + params={"include_source": True}, + ) + + assert response.status_code == 200 + assert response.json()["source"] == {"status": "redacted", "fidelity": "redacted"} + + @pytest.mark.usefixtures("make_dag_with_multiple_versions") + @mock.patch("airflow.api_fastapi.core_api.routes.public.dag_versions.get_auth_manager", autospec=True) + def test_get_dag_version_diff_redacts_source_without_code_access( + self, mock_get_auth_manager, test_client + ): + mock_get_auth_manager.return_value.is_authorized_dag.return_value = False + executed_statements = [] + + def capture_statement(_conn, _cursor, statement, _parameters, _context, _executemany): + executed_statements.append(statement.upper()) + + event.listen(settings.engine, "before_cursor_execute", capture_statement) + try: + response = test_client.get( + "/dags/dag_with_multiple_versions/dagVersions/1/diff/3", + params={"include_source": True}, + ) + finally: + event.remove(settings.engine, "before_cursor_execute", capture_statement) + + assert response.status_code == 200 + assert response.json()["source"] == {"status": "redacted", "fidelity": "redacted"} + assert all("DAG_CODE" not in statement for statement in executed_statements) + + @pytest.mark.usefixtures("make_dag_with_multiple_versions") + @mock.patch( + "airflow.api_fastapi.core_api.routes.public.dag_versions._get_bundle_team_names", + autospec=True, + ) + @mock.patch("airflow.api_fastapi.core_api.routes.public.dag_versions.get_auth_manager", autospec=True) + def test_get_dag_version_diff_authorizes_raw_data_against_historical_bundle( + self, + mock_get_auth_manager, + mock_get_bundle_team_names, + test_client, + session, + ): + base_version = session.scalar( + select(DagVersion).where( + DagVersion.dag_id == "dag_with_multiple_versions", + DagVersion.version_number == 1, + ) + ) + base_version.bundle_name = "another_bundle_name" + session.commit() + mock_get_bundle_team_names.return_value = { + "another_bundle_name": "old-team", + "dag_maker": "new-team", + } + auth_manager = mock_get_auth_manager.return_value + auth_manager.is_authorized_dag.side_effect = lambda **kwargs: ( + kwargs["details"].team_name == "new-team" + ) + + response = test_client.get( + "/dags/dag_with_multiple_versions/dagVersions/1/diff/3", + params={"include_values": True, "include_source": True}, + ) + + assert response.status_code == 200 + response_data = response.json() + assert response_data["values"] == {"status": "unavailable"} + assert response_data["source"] == {"status": "redacted", "fidelity": "redacted"} + assert all( + "before_value" not in change and "after_value" not in change + for change in response_data["changes"] + ) + assert [ + call.kwargs["details"].team_name for call in auth_manager.is_authorized_dag.call_args_list + ] == ["old-team"] + + @pytest.mark.usefixtures("make_dag_with_multiple_versions") + def test_get_dag_version_diff_404(self, test_client): + response = test_client.get("/dags/dag_with_multiple_versions/dagVersions/1/diff/99") + + assert response.status_code == 404 + assert response.json() == { + "detail": "The DagVersion with dag_id: `dag_with_multiple_versions` and version_number: `99` was not found", + } + + def test_should_respond_401(self, unauthenticated_test_client): + response = unauthenticated_test_client.get("/dags/dag_with_multiple_versions/dagVersions/1/diff/2") + + assert response.status_code == 401 + + def test_should_respond_403(self, unauthorized_test_client): + response = unauthorized_test_client.get("/dags/dag_with_multiple_versions/dagVersions/1/diff/2") + + assert response.status_code == 403 + + class TestGetDagVersions(TestDagVersionEndpoint): @pytest.mark.parametrize( ("dag_id", "expected_response", "expected_query_count"), diff --git a/airflow-core/tests/unit/cli/commands/test_dag_command.py b/airflow-core/tests/unit/cli/commands/test_dag_command.py index 059ce1d2f2608..f01b735eb2132 100644 --- a/airflow-core/tests/unit/cli/commands/test_dag_command.py +++ b/airflow-core/tests/unit/cli/commands/test_dag_command.py @@ -30,6 +30,8 @@ import pytest import time_machine from sqlalchemy import func, select +from sqlalchemy.engine import ScalarResult +from sqlalchemy.orm import Session from airflow import settings from airflow._shared.timezones import timezone @@ -80,6 +82,264 @@ dec_27 = DEFAULT_DATE + timedelta(days=-5) +class TestDagVersionDiffCommand: + @pytest.fixture + def parser(self): + return cli_parser.get_parser() + + @mock.patch("airflow.cli.commands.dag_command._print_dag_version_diff", autospec=True) + @mock.patch("airflow.cli.commands.dag_command.get_dag_version_diff", autospec=True) + def test_explicit_versions(self, mock_get_dag_version_diff, mock_print, parser): + result = {"mode": "observed_state"} + mock_get_dag_version_diff.return_value = result + session = MagicMock(spec=Session) + args = parser.parse_args( + [ + "dags", + "versions", + "diff", + "example", + "--from-version", + "2", + "--to-version", + "4", + "--include-values", + "--include-source", + "--max-changes", + "25", + "--output", + "json", + ] + ) + + dag_command.dag_version_diff(args, session=session) + + mock_get_dag_version_diff.assert_called_once_with( + "example", + 2, + 4, + include_values=True, + include_source=True, + max_changes=25, + session=session, + ) + mock_print.assert_called_once_with( + result, + dag_id="example", + base_version_number=2, + target_version_number=4, + output="json", + ) + + @mock.patch("airflow.cli.commands.dag_command._print_dag_version_diff", autospec=True) + @mock.patch("airflow.cli.commands.dag_command.get_dag_version_diff", autospec=True) + def test_previous_to_latest(self, mock_get_dag_version_diff, mock_print, parser): + result = {"mode": "observed_state"} + mock_get_dag_version_diff.return_value = result + scalar_result = MagicMock(spec=ScalarResult) + scalar_result.all.return_value = [5, 3] + session = MagicMock(spec=Session) + session.scalars.return_value = scalar_result + args = parser.parse_args(["dags", "versions", "diff", "example", "--previous-to-latest"]) + + dag_command.dag_version_diff(args, session=session) + + mock_get_dag_version_diff.assert_called_once_with( + "example", + 3, + 5, + include_values=False, + include_source=False, + max_changes=500, + session=session, + ) + mock_print.assert_called_once_with( + result, + dag_id="example", + base_version_number=3, + target_version_number=5, + output="table", + ) + + @pytest.mark.parametrize( + ("options", "message"), + [ + (["--from-version", "1"], "Provide both --from-version and --to-version"), + ( + ["--previous-to-latest", "--from-version", "1", "--to-version", "2"], + "--previous-to-latest cannot be combined", + ), + ], + ) + def test_rejects_invalid_version_options(self, parser, options, message): + session = MagicMock(spec=Session) + args = parser.parse_args(["dags", "versions", "diff", "example", *options]) + + with pytest.raises(SystemExit, match=message): + dag_command.dag_version_diff(args, session=session) + + def test_previous_to_latest_requires_two_versions(self, parser): + scalar_result = MagicMock(spec=ScalarResult) + scalar_result.all.return_value = [1] + session = MagicMock(spec=Session) + session.scalars.return_value = scalar_result + args = parser.parse_args(["dags", "versions", "diff", "example", "--previous-to-latest"]) + + with pytest.raises(SystemExit, match="must have at least two versions"): + dag_command.dag_version_diff(args, session=session) + + @pytest.mark.parametrize( + "error", + [ + dag_command.DagVersionNotFoundError("missing version"), + ValueError("invalid maximum"), + ], + ) + @mock.patch("airflow.cli.commands.dag_command.get_dag_version_diff", autospec=True) + def test_translates_service_errors(self, mock_get_dag_version_diff, parser, error): + mock_get_dag_version_diff.side_effect = error + session = MagicMock(spec=Session) + args = parser.parse_args( + [ + "dags", + "versions", + "diff", + "example", + "--from-version", + "1", + "--to-version", + "2", + ] + ) + + with pytest.raises(SystemExit, match=str(error)): + dag_command.dag_version_diff(args, session=session) + + @pytest.mark.parametrize("output", ["table", "plain"]) + def test_human_output_includes_values_and_truncation(self, stdout_capture, output): + result = { + "mode": "observed_state", + "changes": [ + { + "path": "/dag/value", + "operation": "changed", + "category": "unknown", + "impact": "unknown", + "before_value": 1, + "after_value": None, + }, + { + "path": "/dag/added", + "operation": "added", + "category": "unknown", + "impact": "unknown", + "after_value": "new", + }, + ], + "source": {"status": "unavailable", "fidelity": "unavailable"}, + "values": {"status": "available"}, + "truncated": True, + } + + with stdout_capture as stdout: + dag_command._print_dag_version_diff( + result, + dag_id="example", + base_version_number=1, + target_version_number=2, + output=output, + ) + + rendered = stdout.getvalue() + assert "before_value" in rendered + assert "after_value" in rendered + assert "null" in rendered + assert "" in rendered + assert "new" in rendered + assert "Changes truncated" in rendered + + def test_human_output_reports_unavailable_reason(self, stdout_capture): + result = { + "mode": "unavailable", + "changes": [], + "source": {"status": "unavailable", "fidelity": "unavailable"}, + "truncated": False, + "unavailable_reason": "serialized_dag_missing", + } + + with stdout_capture as stdout: + dag_command._print_dag_version_diff( + result, + dag_id="example", + base_version_number=1, + target_version_number=2, + output="plain", + ) + + rendered = stdout.getvalue() + assert "No changes" in rendered + assert "Source: unavailable" in rendered + assert "Unavailable reason: serialized_dag_missing" in rendered + + def test_human_output_without_values_includes_source(self, stdout_capture): + result = { + "mode": "observed_state", + "changes": [ + { + "path": "/dag/value", + "operation": "changed", + "category": "unknown", + "impact": "unknown", + } + ], + "source": { + "status": "current_stored_code", + "fidelity": "current_stored_code", + "base": {"content": "base source"}, + "target": {"content": "target source"}, + }, + "truncated": False, + } + + with stdout_capture as stdout: + dag_command._print_dag_version_diff( + result, + dag_id="example", + base_version_number=1, + target_version_number=2, + output="plain", + ) + + rendered = stdout.getvalue() + assert "before_value" not in rendered + assert "after_value" not in rendered + assert "Base source:\nbase source" in rendered + assert "Target source:\ntarget source" in rendered + + @pytest.mark.parametrize( + ("output", "renderer_name"), + [("json", "print_as_json"), ("yaml", "print_as_yaml")], + ) + @mock.patch("airflow.cli.commands.dag_command.AirflowConsole", autospec=True) + def test_structured_output(self, mock_console, output, renderer_name): + result = { + "mode": "observed_state", + "changes": [], + "source": {"status": "unavailable", "fidelity": "unavailable"}, + "truncated": False, + } + + dag_command._print_dag_version_diff( + result, + dag_id="example", + base_version_number=1, + target_version_number=2, + output=output, + ) + + getattr(mock_console.return_value, renderer_name).assert_called_once_with(result) + + class TestCliDags: parser: argparse.ArgumentParser diff --git a/airflow-core/tests/unit/cli/test_cli_parser.py b/airflow-core/tests/unit/cli/test_cli_parser.py index acba2fa5d2e26..b65e7c3797e38 100644 --- a/airflow-core/tests/unit/cli/test_cli_parser.py +++ b/airflow-core/tests/unit/cli/test_cli_parser.py @@ -55,6 +55,34 @@ class TestCli: + def test_dag_version_diff_parser(self): + args = cli_parser.get_parser().parse_args( + [ + "dags", + "versions", + "diff", + "example", + "--from-version", + "12", + "--to-version", + "13", + "--include-values", + "--include-source", + "--max-changes", + "25", + "--output", + "json", + ] + ) + + assert args.dag_id == "example" + assert args.from_version == 12 + assert args.to_version == 13 + assert args.include_values is True + assert args.include_source is True + assert args.max_changes == 25 + assert args.output == "json" + def test_arg_option_long_only(self): """ Test if the name of cli.args long option valid diff --git a/airflow-core/tests/unit/models/test_dag_version_diff.py b/airflow-core/tests/unit/models/test_dag_version_diff.py new file mode 100644 index 0000000000000..b072ca5782b31 --- /dev/null +++ b/airflow-core/tests/unit/models/test_dag_version_diff.py @@ -0,0 +1,349 @@ +# 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. + +from __future__ import annotations + +import pytest + +from airflow.models.dag_version_diff import build_serialized_dag_diff + + +def _payload( + *, + tasks: list[dict], + tags: list[str] | None = None, + schedule: str = "daily", + dependencies: list[dict] | None = None, +) -> dict: + return { + "__version": 3, + "dag": { + "dag_id": "example", + "schedule": schedule, + "tags": tags or [], + "tasks": [ + { + "__type": "airflow.providers.standard.operators.empty.EmptyOperator", + "__var": task, + } + for task in tasks + ], + "dag_dependencies": dependencies or [], + }, + } + + +def test_build_diff_is_deterministic_and_normalizes_order() -> None: + base = _payload( + tasks=[{"task_id": "extract", "retries": 1}, {"task_id": "load", "retries": 1}], + tags=["one", "two"], + dependencies=[ + { + "dependency_type": "task", + "dependency_id": "extract-load", + "source": "extract", + "target": "load", + "label": "extract-load", + } + ], + ) + target = _payload( + tasks=[{"task_id": "load", "retries": 1}, {"task_id": "extract", "retries": 1}], + tags=["two", "one"], + dependencies=[ + { + "label": "extract-load", + "target": "load", + "source": "extract", + "dependency_id": "extract-load", + "dependency_type": "task", + } + ], + ) + + result = build_serialized_dag_diff(base_data=base, target_data=target) + + assert result["mode"] == "observed_state" + assert result["serialized_dag_schema_versions"] == {"base": 3, "target": 3} + assert result["changes"] == [] + assert result["truncated"] is False + + +def test_build_diff_reports_categories_digests_and_values() -> None: + base = _payload(tasks=[{"task_id": "extract", "retries": 1}], tags=["old"]) + target = _payload( + tasks=[ + {"task_id": "extract", "retries": 2}, + {"task_id": "load", "retries": 1}, + ], + tags=["new"], + schedule="hourly", + ) + + result = build_serialized_dag_diff( + base_data=base, + target_data=target, + base_provenance={"bundle_version": "one"}, + target_provenance={"bundle_version": "two"}, + include_values=True, + ) + + changes = {change["path"]: change for change in result["changes"]} + assert changes["/dag/tasks/extract/retries"]["category"] == "task" + assert changes["/dag/tasks/extract/retries"]["impact"] == "execution" + assert changes["/dag/tasks/extract/retries"]["before_value"] == 1 + assert changes["/dag/tasks/extract/retries"]["after_value"] == 2 + assert changes["/dag/tasks/extract/retries"]["before_digest"].startswith("sha256:") + assert changes["/dag/tasks/load"]["operation"] == "added" + assert changes["/dag/schedule"]["category"] == "schedule" + assert changes["/dag/tags/old"]["operation"] == "removed" + assert changes["/dag/tags/new"]["operation"] == "added" + assert changes["/dag/tags/old"]["category"] == "metadata" + assert changes["/provenance/bundle_version"]["category"] == "provenance" + assert changes["/provenance/bundle_version"]["impact"] == "provenance" + + +def test_build_diff_bounds_changes_and_reports_truncation() -> None: + result = build_serialized_dag_diff( + base_data=_payload(tasks=[{"task_id": "extract", "retries": 1}], tags=["old"]), + target_data=_payload(tasks=[{"task_id": "extract", "retries": 2}], tags=["new"]), + max_changes=1, + ) + + assert len(result["changes"]) == 1 + assert result["truncated"] is True + + +def test_build_diff_reports_dependency_changes() -> None: + dependency = { + "dependency_type": "sensor", + "dependency_id": "upstream-task", + "source": "upstream-dag", + "target": "example", + "label": "upstream-task", + } + + result = build_serialized_dag_diff( + base_data=_payload(tasks=[], dependencies=[dependency]), + target_data=_payload(tasks=[]), + ) + + assert len(result["changes"]) == 1 + assert result["changes"][0]["category"] == "dependency" + assert result["changes"][0]["impact"] == "execution" + + +def test_build_diff_classifies_fail_fast_as_schedule() -> None: + base = _payload(tasks=[]) + target = _payload(tasks=[]) + base["dag"]["fail_fast"] = False + target["dag"]["fail_fast"] = True + + result = build_serialized_dag_diff(base_data=base, target_data=target) + + assert len(result["changes"]) == 1 + change = result["changes"][0] + assert change["path"] == "/dag/fail_fast" + assert change["operation"] == "changed" + assert change["category"] == "schedule" + assert change["impact"] == "execution" + + +@pytest.mark.parametrize( + ("base_data", "target_data", "reason"), + [ + (None, _payload(tasks=[]), "serialized_dag_missing"), + ({"__version": 99, "dag": {}}, _payload(tasks=[]), "unsupported_serialized_dag_schema_version:99"), + ({"__version": 3, "dag": []}, _payload(tasks=[]), "serialized_dag_canonicalization_failed"), + ( + { + "__version": 1, + "dag": { + "dag_id": "example", + "tasks": [], + "task_group": {}, + "schedule_interval": 1, + }, + }, + _payload(tasks=[]), + "serialized_dag_canonicalization_failed", + ), + ( + { + "__version": 1, + "dag": { + "dag_id": "example", + "tasks": [], + "task_group": {}, + "schedule_interval": {"__type": "timedelta", "__var": 10**20}, + }, + }, + _payload(tasks=[]), + "serialized_dag_canonicalization_failed", + ), + ], +) +def test_build_diff_returns_unavailable_for_unsafe_inputs(base_data, target_data, reason) -> None: + result = build_serialized_dag_diff(base_data=base_data, target_data=target_data) + + assert result["mode"] == "unavailable" + assert result["changes"] == [] + assert result["unavailable_reason"] == reason + + +def test_build_diff_rejects_non_positive_change_bound() -> None: + with pytest.raises(ValueError, match="max_changes must be a positive integer"): + build_serialized_dag_diff(base_data=_payload(tasks=[]), target_data=_payload(tasks=[]), max_changes=0) + + +def test_build_diff_rejects_unbounded_change_bound() -> None: + with pytest.raises(ValueError, match="max_changes must not exceed 5000"): + build_serialized_dag_diff( + base_data=_payload(tasks=[]), target_data=_payload(tasks=[]), max_changes=5001 + ) + + +def test_build_diff_reports_unkeyed_lists_as_one_stable_change() -> None: + base = _payload(tasks=[]) + target = _payload(tasks=[]) + base["dag"]["custom_list"] = [{"name": "old"}] + target["dag"]["custom_list"] = [{"name": "new"}] + + result = build_serialized_dag_diff(base_data=base, target_data=target) + + assert len(result["changes"]) == 1 + change = result["changes"][0] + assert change["path"] == "/dag/custom_list" + assert change["operation"] == "changed" + assert change["category"] == "unknown" + assert change["impact"] == "unknown" + assert change["before_digest"].startswith("sha256:") + assert change["after_digest"].startswith("sha256:") + + +def test_build_diff_reports_deadline_lists_as_one_stable_change() -> None: + base = _payload(tasks=[]) + target = _payload(tasks=[]) + base["dag"]["deadline"] = [{"name": "old", "interval": 60}] + target["dag"]["deadline"] = [{"name": "new", "interval": 60}] + + result = build_serialized_dag_diff(base_data=base, target_data=target) + + assert result["mode"] == "observed_state" + assert len(result["changes"]) == 1 + change = result["changes"][0] + assert change["path"] == "/dag/deadline" + assert change["operation"] == "changed" + assert change["category"] == "deadline" + assert change["impact"] == "execution" + + +def test_build_diff_preserves_order_sensitive_string_lists() -> None: + base = _payload(tasks=[]) + target = _payload(tasks=[]) + base["dag"]["template_searchpath"] = ["first", "second"] + target["dag"]["template_searchpath"] = ["second", "first"] + + result = build_serialized_dag_diff(base_data=base, target_data=target) + + assert len(result["changes"]) == 1 + change = result["changes"][0] + assert change["path"] == "/dag/template_searchpath" + assert change["operation"] == "changed" + + +@pytest.mark.parametrize( + ("base_task", "target_task"), + [ + ({"task_id": "extract", "retries": 2}, {"task_id": "extract"}), + ( + {"task_id": "extract", "retries": 2, "partial_kwargs": {"retries": 2}}, + {"task_id": "extract", "partial_kwargs": {}}, + ), + ], +) +def test_build_diff_normalizes_v3_client_defaults(base_task, target_task) -> None: + base = _payload(tasks=[base_task]) + base["__version"] = 2 + target = _payload(tasks=[target_task]) + target["client_defaults"] = {"tasks": {"retries": 2}} + + result = build_serialized_dag_diff(base_data=base, target_data=target) + + assert result["serialized_dag_schema_versions"] == {"base": 2, "target": 3} + assert result["mode"] == "observed_state" + assert result["changes"] == [] + + +@pytest.mark.parametrize( + "base_data", + [ + {"__version": 3, "dag": {"tasks": []}, "client_defaults": []}, + {"__version": 3, "dag": {"tasks": []}, "client_defaults": {"tasks": []}}, + {"__version": 3, "dag": {"tasks": {}}, "client_defaults": {"tasks": {}}}, + {"__version": 3, "dag": {"tasks": [{}]}, "client_defaults": {"tasks": {}}}, + ], +) +def test_build_diff_rejects_malformed_client_defaults(base_data) -> None: + result = build_serialized_dag_diff(base_data=base_data, target_data=_payload(tasks=[])) + + assert result["mode"] == "unavailable" + assert result["changes"] == [] + assert result["unavailable_reason"] == "serialized_dag_canonicalization_failed" + + +def test_build_diff_ignores_exact_duplicate_dependencies() -> None: + dependency = { + "dependency_type": "task", + "dependency_id": "extract-load", + "source": "extract", + "target": "load", + "label": "extract-load", + } + + result = build_serialized_dag_diff( + base_data=_payload(tasks=[], dependencies=[dependency, dependency]), + target_data=_payload(tasks=[], dependencies=[dependency]), + ) + + assert result["mode"] == "observed_state" + assert result["changes"] == [] + + +def test_build_diff_dependency_keys_do_not_collide_on_delimiters() -> None: + first_dependency = { + "dependency_type": "trigger", + "dependency_id": "id", + "source": "a", + "target": "b|c", + "label": "d", + } + second_dependency = { + "dependency_type": "trigger", + "dependency_id": "id", + "source": "a", + "target": "b", + "label": "c|d", + } + + result = build_serialized_dag_diff( + base_data=_payload(tasks=[], dependencies=[first_dependency, second_dependency]), + target_data=_payload(tasks=[], dependencies=[second_dependency, first_dependency]), + ) + + assert result["mode"] == "observed_state" + assert result["changes"] == [] diff --git a/airflow-ctl/src/airflowctl/api/datamodels/generated.py b/airflow-ctl/src/airflowctl/api/datamodels/generated.py index 9bb75d8dd87ce..54fb5e29fcdfc 100644 --- a/airflow-ctl/src/airflowctl/api/datamodels/generated.py +++ b/airflow-ctl/src/airflowctl/api/datamodels/generated.py @@ -622,6 +622,86 @@ class DagTagResponse(BaseModel): dag_display_name: Annotated[str, Field(title="Dag Display Name")] +class Operation(str, Enum): + ADDED = "added" + REMOVED = "removed" + CHANGED = "changed" + + +class Category(str, Enum): + TASK = "task" + DEPENDENCY = "dependency" + SCHEDULE = "schedule" + PARAM = "param" + ASSET = "asset" + CALLBACK = "callback" + DEADLINE = "deadline" + METADATA = "metadata" + PROVENANCE = "provenance" + UNKNOWN = "unknown" + + +class Impact(str, Enum): + EXECUTION = "execution" + METADATA = "metadata" + PROVENANCE = "provenance" + UNKNOWN = "unknown" + + +class DagVersionDiffChange(BaseModel): + """ + One observed-state change between two serialized Dag versions. + """ + + path: Annotated[str, Field(title="Path")] + operation: Annotated[Operation, Field(title="Operation")] + category: Annotated[Category, Field(title="Category")] + impact: Annotated[Impact, Field(title="Impact")] + before_digest: Annotated[str | None, Field(title="Before Digest")] = None + after_digest: Annotated[str | None, Field(title="After Digest")] = None + before_value: Annotated[Any, Field(title="Before Value")] = None + after_value: Annotated[Any, Field(title="After Value")] = None + + +class Mode(str, Enum): + OBSERVED_STATE = "observed_state" + UNAVAILABLE = "unavailable" + + +class Status(str, Enum): + CURRENT_STORED_CODE = "current_stored_code" + REDACTED = "redacted" + UNAVAILABLE = "unavailable" + + +class Fidelity(str, Enum): + CURRENT_STORED_CODE = "current_stored_code" + REDACTED = "redacted" + UNAVAILABLE = "unavailable" + + +class DagVersionDiffSourceSide(BaseModel): + """ + Source metadata for one side of a Dag version comparison. + """ + + digest: Annotated[str | None, Field(title="Digest")] = None + content: Annotated[str | None, Field(title="Content")] = None + + +class Status1(str, Enum): + AVAILABLE = "available" + UNAVAILABLE = "unavailable" + + +class DagVersionDiffValues(BaseModel): + """ + Visibility metadata for raw serialized Dag values. + """ + + status: Annotated[Status1, Field(title="Status")] + + class DagVersionResponse(BaseModel): """ Dag Version serializer for responses. @@ -1907,6 +1987,18 @@ class DagStatsResponse(BaseModel): stats: Annotated[list[DagStatsStateResponse], Field(title="Stats")] +class DagVersionDiffSource(BaseModel): + """ + Source comparison metadata. + """ + + status: Annotated[Status, Field(title="Status")] + fidelity: Annotated[Fidelity, Field(title="Fidelity")] + changed: Annotated[bool | None, Field(title="Changed")] = None + base: DagVersionDiffSourceSide | None = None + target: DagVersionDiffSourceSide | None = None + + class DryRunBackfillCollectionResponse(BaseModel): """ Backfill collection serializer for responses in dry-run mode. @@ -2375,6 +2467,23 @@ class DagStatsCollectionResponse(BaseModel): total_entries: Annotated[int, Field(title="Total Entries")] +class DagVersionDiffResponse(BaseModel): + """ + Observed-state diff response for two Dag versions. + """ + + diff_schema_version: Annotated[int, Field(title="Diff Schema Version")] + serialized_dag_schema_versions: Annotated[ + dict[str, int | None], Field(title="Serialized Dag Schema Versions") + ] + mode: Annotated[Mode, Field(title="Mode")] + changes: Annotated[list[DagVersionDiffChange], Field(title="Changes")] + source: DagVersionDiffSource + values: DagVersionDiffValues | None = None + truncated: Annotated[bool, Field(title="Truncated")] + unavailable_reason: Annotated[str | None, Field(title="Unavailable Reason")] = None + + class HITLDetail(BaseModel): """ Schema for Human-in-the-loop detail.