Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/execution_api/versions/v2026_06_30.py: 57%
64 statements
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 14:22 +0000
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-07 14:22 +0000
1# Licensed to the Apache Software Foundation (ASF) under one
2# or more contributor license agreements. See the NOTICE file
3# distributed with this work for additional information
4# regarding copyright ownership. The ASF licenses this file
5# to you under the Apache License, Version 2.0 (the
6# "License"); you may not use this file except in compliance
7# with the License. You may obtain a copy of the License at
8#
9# http://www.apache.org/licenses/LICENSE-2.0
10#
11# Unless required by applicable law or agreed to in writing,
12# software distributed under the License is distributed on an
13# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14# KIND, either express or implied. See the License for the
15# specific language governing permissions and limitations
16# under the License.
18from __future__ import annotations
20from cadwyn import (
21 ResponseInfo,
22 VersionChange,
23 convert_response_to_previous_version_for,
24 endpoint,
25 schema,
26)
28from airflow.api_fastapi.execution_api.datamodels.taskinstance import (
29 AssetEventDagRunReference,
30 DagRun,
31 TaskInstance,
32 TIAwaitingInputStatePayload,
33 TIRetryStatePayload,
34 TIRunContext,
35)
38class AddVariableKeysEndpoint(VersionChange):
39 """Add GET /variables/keys endpoint for listing variable keys with optional prefix filter."""
41 description = __doc__
43 instructions_to_migrate_to_previous_version = (endpoint("/variables/keys", ["GET"]).didnt_exist,)
46class AddConnectionTestEndpoint(VersionChange):
47 """Add connection-tests endpoints for the async connection-test workflow."""
49 description = __doc__
51 instructions_to_migrate_to_previous_version = (
52 endpoint("/connection-tests/{connection_test_id}", ["PATCH"]).didnt_exist,
53 endpoint("/connection-tests/{connection_test_id}/connection", ["GET"]).didnt_exist,
54 )
57class AddTaskInstanceQueueField(VersionChange):
58 """Add the `queue` field to the TaskInstance model."""
60 description = __doc__
62 instructions_to_migrate_to_previous_version = (schema(TaskInstance).field("queue").didnt_exist,)
65class AddAwaitingInputStatePayload(VersionChange):
66 """Add the awaiting_input task instance state transition payload (Human-in-the-loop, no trigger)."""
68 description = __doc__
70 instructions_to_migrate_to_previous_version = (
71 schema(TIAwaitingInputStatePayload).field("state").didnt_exist,
72 schema(TIAwaitingInputStatePayload).field("timeout").didnt_exist,
73 schema(TIAwaitingInputStatePayload).field("next_method").didnt_exist,
74 schema(TIAwaitingInputStatePayload).field("next_kwargs").didnt_exist,
75 schema(TIAwaitingInputStatePayload).field("rendered_map_index").didnt_exist,
76 )
79class AddRetryPolicyFields(VersionChange):
80 """Add retry_delay_seconds and retry_reason fields to TIRetryStatePayload for pluggable retry policies."""
82 description = __doc__
84 instructions_to_migrate_to_previous_version = (
85 schema(TIRetryStatePayload).field("retry_delay_seconds").didnt_exist,
86 schema(TIRetryStatePayload).field("retry_reason").didnt_exist,
87 )
90class AddTeamNameField(VersionChange):
91 """Add the ``team_name`` field to DagRun model."""
93 description = __doc__
95 instructions_to_migrate_to_previous_version = (schema(DagRun).field("team_name").didnt_exist,)
97 @convert_response_to_previous_version_for(TIRunContext) # type: ignore[arg-type]
98 def remove_team_name_field(response: ResponseInfo) -> None: # type: ignore[misc]
99 """Remove the ``team_name`` field from dag_run for older API versions."""
100 if "dag_run" in response.body and isinstance(response.body["dag_run"], dict):
101 response.body["dag_run"].pop("team_name", None)
103 @convert_response_to_previous_version_for(DagRun) # type: ignore[arg-type]
104 def remove_team_name_from_dag_run_response(response: ResponseInfo) -> None: # type: ignore[misc]
105 """Remove the ``team_name`` field from responses returning a DagRun directly."""
106 if isinstance(response.body, dict):
107 response.body.pop("team_name", None)
109 # Schema-based converters are matched against the route's response model by identity, so a
110 # route annotated ``DagRun | None`` never matches ``DagRun`` and has to be addressed by path.
111 @convert_response_to_previous_version_for("/dag-runs/previous", ["GET"]) # type: ignore[arg-type]
112 def remove_team_name_from_previous_dag_run(response: ResponseInfo) -> None: # type: ignore[misc]
113 """Remove the ``team_name`` field from the previous-run response."""
114 if isinstance(response.body, dict):
115 response.body.pop("team_name", None)
118class AddAssetsByAliasEndpoint(VersionChange):
119 """Add endpoint to resolve assets from an AssetAlias."""
121 description = __doc__
123 instructions_to_migrate_to_previous_version = (endpoint("/assets/by-alias", ["GET"]).didnt_exist,)
126class AddTaskAndAssetStateStoreEndpoints(VersionChange):
127 """Add task state store and asset state store API endpoints."""
129 description = __doc__
131 instructions_to_migrate_to_previous_version = (
132 endpoint("/store/ti/{task_instance_id}/{key:path}", ["GET"]).didnt_exist,
133 endpoint("/store/ti/{task_instance_id}/{key:path}", ["PUT"]).didnt_exist,
134 endpoint("/store/ti/{task_instance_id}/{key:path}", ["DELETE"]).didnt_exist,
135 endpoint("/store/ti/{task_instance_id}", ["DELETE"]).didnt_exist,
136 endpoint("/store/asset/by-name/value", ["GET"]).didnt_exist,
137 endpoint("/store/asset/by-name/value", ["PUT"]).didnt_exist,
138 endpoint("/store/asset/by-name/value", ["DELETE"]).didnt_exist,
139 endpoint("/store/asset/by-name/clear", ["DELETE"]).didnt_exist,
140 endpoint("/store/asset/by-uri/value", ["GET"]).didnt_exist,
141 endpoint("/store/asset/by-uri/value", ["PUT"]).didnt_exist,
142 endpoint("/store/asset/by-uri/value", ["DELETE"]).didnt_exist,
143 endpoint("/store/asset/by-uri/clear", ["DELETE"]).didnt_exist,
144 )
147class AddPartitionDateField(VersionChange):
148 """Expose the consumer DagRun's partition datetime on the execution API so consumer tasks can template it."""
150 description = __doc__
152 instructions_to_migrate_to_previous_version = (schema(DagRun).field("partition_date").didnt_exist,)
154 @convert_response_to_previous_version_for(TIRunContext) # type: ignore[arg-type]
155 def remove_partition_date_from_dag_run(response: ResponseInfo) -> None: # type: ignore[misc]
156 """Strip ``partition_date`` from the nested ``dag_run`` payload for older clients."""
157 if "dag_run" in response.body and isinstance(response.body["dag_run"], dict):
158 response.body["dag_run"].pop("partition_date", None)
160 @convert_response_to_previous_version_for(DagRun) # type: ignore[arg-type]
161 def remove_partition_date_from_dag_run_response(response: ResponseInfo) -> None: # type: ignore[misc]
162 """Strip ``partition_date`` from responses returning a DagRun directly."""
163 if isinstance(response.body, dict):
164 response.body.pop("partition_date", None)
166 # See remove_team_name_from_previous_dag_run: ``DagRun | None`` needs a path-based converter.
167 @convert_response_to_previous_version_for("/dag-runs/previous", ["GET"]) # type: ignore[arg-type]
168 def remove_partition_date_from_previous_dag_run(response: ResponseInfo) -> None: # type: ignore[misc]
169 """Strip ``partition_date`` from the previous-run response."""
170 if isinstance(response.body, dict):
171 response.body.pop("partition_date", None)
174class AddConsumedAssetEventPartitionKeyField(VersionChange):
175 """Expose the upstream partition key on the asset events that triggered a consumer Dag run."""
177 description = __doc__
179 instructions_to_migrate_to_previous_version = (
180 schema(AssetEventDagRunReference).field("partition_key").didnt_exist,
181 )
183 # Only ``TIRunContext`` can carry these events: ``DagRun.safe_extract_from_orm`` defaults
184 # ``consumed_asset_events`` to ``[]`` whenever the relationship is not already loaded, and
185 # ``/run`` is the only route that loads it -- so bare ``DagRun`` responses never carry one.
186 @convert_response_to_previous_version_for(TIRunContext) # type: ignore[arg-type]
187 def remove_partition_key_from_consumed_asset_events(response: ResponseInfo) -> None: # type: ignore[misc]
188 """Strip ``partition_key`` from each consumed asset event for older clients."""
189 dag_run = response.body.get("dag_run")
190 if isinstance(dag_run, dict):
191 for event in dag_run.get("consumed_asset_events") or ():
192 if isinstance(event, dict):
193 event.pop("partition_key", None)