Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/execution_api/versions/v2025_05_20.py: 31%

18 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 ResponseInfo, VersionChange, convert_response_to_previous_version_for, schema 

21 

22from airflow.api_fastapi.execution_api.datamodels.taskinstance import TIRunContext 

23 

24 

25class DowngradeUpstreamMapIndexes(VersionChange): 

26 """Downgrade the upstream map indexes type for older clients.""" 

27 

28 description = __doc__ 

29 

30 instructions_to_migrate_to_previous_version = ( 

31 schema(TIRunContext).field("upstream_map_indexes").had(type=dict[str, int | None] | None), 

32 ) 

33 

34 @convert_response_to_previous_version_for(TIRunContext) # type: ignore[arg-type] 

35 def downgrade_upstream_map_indexes(response: ResponseInfo = None) -> None: # type: ignore 

36 """ 

37 Downgrades the `upstream_map_indexes` field when converting to the previous version. 

38 

39 Ensures that the field is only a dictionary of [str, int] (old format). 

40 """ 

41 resp = response.body.get("upstream_map_indexes") 

42 if isinstance(resp, dict): 

43 downgraded: dict[str, int | list | None] = {} 

44 for k, v in resp.items(): 

45 if isinstance(v, int): 

46 downgraded[k] = v 

47 elif isinstance(v, list) and v and all(isinstance(i, int) for i in v): 

48 downgraded[k] = v[0] 

49 else: 

50 # Keep values like None as is — the Task SDK expects them unchanged during mapped task expansion, 

51 # and modifying them can cause unexpected failures. 

52 downgraded[k] = None 

53 response.body["upstream_map_indexes"] = downgraded