From ec383a640d0e6e35d4d0f7d825aae97271524daa Mon Sep 17 00:00:00 2001 From: Ephraim Anierobi Date: Mon, 13 Jul 2026 11:31:43 +0100 Subject: [PATCH] Add Dag version diff tooling Dag authors and operators need a reliable way to understand how two stored Dag versions differ without downloading serialized payloads or manually comparing implementation details. This provides a semantic comparison surface shared by the public API and the airflow dags version-diff command, with machine-readable and human-readable output for automation and investigation. Historical serialized Dags can use different schema shapes, so comparisons normalize task, schedule, dependency, and metadata state before calculating changes. Optional source and rendered-value data are explicitly classified as available, redacted, or unavailable, and source authorization is evaluated against each version's historical bundle to avoid exposing code across team boundaries. Malformed legacy payloads and missing historical data degrade to unavailable results instead of failing the whole request. The OpenAPI contract and regression coverage keep authorization, duplicate dependency handling, schema conversion, CLI parsing and output, and failure states aligned across the feature. --- .../docs/security/api_permissions_ref.rst | 4 + .../core_api/datamodels/dag_versions.py | 61 +++ .../openapi/v2-rest-api-generated.yaml | 257 ++++++++++ .../core_api/routes/public/dag_versions.py | 200 +++++++- airflow-core/src/airflow/cli/cli_config.py | 61 ++- .../src/airflow/cli/commands/dag_command.py | 111 +++- .../src/airflow/models/dag_version_diff.py | 483 ++++++++++++++++++ .../airflow/ui/openapi-gen/queries/common.ts | 11 + .../ui/openapi-gen/queries/ensureQueryData.ts | 21 + .../ui/openapi-gen/queries/prefetch.ts | 21 + .../airflow/ui/openapi-gen/queries/queries.ts | 21 + .../ui/openapi-gen/queries/suspense.ts | 21 + .../ui/openapi-gen/requests/schemas.gen.ts | 229 +++++++++ .../ui/openapi-gen/requests/services.gen.ts | 38 +- .../ui/openapi-gen/requests/types.gen.ts | 110 +++- .../routes/public/test_dag_versions.py | 267 ++++++++++ .../unit/cli/commands/test_dag_command.py | 260 ++++++++++ .../tests/unit/cli/test_cli_parser.py | 28 + .../unit/models/test_dag_version_diff.py | 349 +++++++++++++ .../airflowctl/api/datamodels/generated.py | 109 ++++ 20 files changed, 2655 insertions(+), 7 deletions(-) create mode 100644 airflow-core/src/airflow/models/dag_version_diff.py create mode 100644 airflow-core/tests/unit/models/test_dag_version_diff.py 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.