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

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. 

17 

18from __future__ import annotations 

19 

20from cadwyn import ( 

21 ResponseInfo, 

22 VersionChange, 

23 convert_response_to_previous_version_for, 

24 endpoint, 

25 schema, 

26) 

27 

28from airflow.api_fastapi.execution_api.datamodels.taskinstance import ( 

29 AssetEventDagRunReference, 

30 DagRun, 

31 TaskInstance, 

32 TIAwaitingInputStatePayload, 

33 TIRetryStatePayload, 

34 TIRunContext, 

35) 

36 

37 

38class AddVariableKeysEndpoint(VersionChange): 

39 """Add GET /variables/keys endpoint for listing variable keys with optional prefix filter.""" 

40 

41 description = __doc__ 

42 

43 instructions_to_migrate_to_previous_version = (endpoint("/variables/keys", ["GET"]).didnt_exist,) 

44 

45 

46class AddConnectionTestEndpoint(VersionChange): 

47 """Add connection-tests endpoints for the async connection-test workflow.""" 

48 

49 description = __doc__ 

50 

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 ) 

55 

56 

57class AddTaskInstanceQueueField(VersionChange): 

58 """Add the `queue` field to the TaskInstance model.""" 

59 

60 description = __doc__ 

61 

62 instructions_to_migrate_to_previous_version = (schema(TaskInstance).field("queue").didnt_exist,) 

63 

64 

65class AddAwaitingInputStatePayload(VersionChange): 

66 """Add the awaiting_input task instance state transition payload (Human-in-the-loop, no trigger).""" 

67 

68 description = __doc__ 

69 

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 ) 

77 

78 

79class AddRetryPolicyFields(VersionChange): 

80 """Add retry_delay_seconds and retry_reason fields to TIRetryStatePayload for pluggable retry policies.""" 

81 

82 description = __doc__ 

83 

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 ) 

88 

89 

90class AddTeamNameField(VersionChange): 

91 """Add the ``team_name`` field to DagRun model.""" 

92 

93 description = __doc__ 

94 

95 instructions_to_migrate_to_previous_version = (schema(DagRun).field("team_name").didnt_exist,) 

96 

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) 

102 

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) 

108 

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) 

116 

117 

118class AddAssetsByAliasEndpoint(VersionChange): 

119 """Add endpoint to resolve assets from an AssetAlias.""" 

120 

121 description = __doc__ 

122 

123 instructions_to_migrate_to_previous_version = (endpoint("/assets/by-alias", ["GET"]).didnt_exist,) 

124 

125 

126class AddTaskAndAssetStateStoreEndpoints(VersionChange): 

127 """Add task state store and asset state store API endpoints.""" 

128 

129 description = __doc__ 

130 

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 ) 

145 

146 

147class AddPartitionDateField(VersionChange): 

148 """Expose the consumer DagRun's partition datetime on the execution API so consumer tasks can template it.""" 

149 

150 description = __doc__ 

151 

152 instructions_to_migrate_to_previous_version = (schema(DagRun).field("partition_date").didnt_exist,) 

153 

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) 

159 

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) 

165 

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) 

172 

173 

174class AddConsumedAssetEventPartitionKeyField(VersionChange): 

175 """Expose the upstream partition key on the asset events that triggered a consumer Dag run.""" 

176 

177 description = __doc__ 

178 

179 instructions_to_migrate_to_previous_version = ( 

180 schema(AssetEventDagRunReference).field("partition_key").didnt_exist, 

181 ) 

182 

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)