Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/datamodels/tasks.py: 84%

77 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 

20import inspect 

21from collections import abc 

22from datetime import datetime 

23from typing import TYPE_CHECKING, Any 

24 

25from pydantic import computed_field, field_validator, model_validator 

26 

27from airflow._shared.module_loading import qualname 

28from airflow.api_fastapi.common.types import TimeDeltaWithValidation 

29from airflow.api_fastapi.core_api.base import BaseModel 

30from airflow.task.priority_strategy import ( 

31 get_weight_rule_from_priority_weight_strategy, 

32 validate_and_load_priority_weight_strategy, 

33) 

34 

35if TYPE_CHECKING: 35 ↛ 36line 35 didn't jump to line 36 because the condition on line 35 was never true

36 from airflow.serialization.definitions.param import SerializedParamsDict 

37 from airflow.task.priority_strategy import PriorityWeightStrategy 

38 

39 

40def _get_class_ref(obj) -> dict[str, str | None]: 

41 """Return the class_ref dict for obj.""" 

42 module_path = getattr(obj, "_task_module", None) 

43 if module_path is None: 43 ↛ 44line 43 didn't jump to line 44 because the condition on line 43 was never true

44 module_type = inspect.getmodule(obj) 

45 module_path = module_type.__name__ if module_type else None 

46 

47 class_name = obj.task_type 

48 

49 return { 

50 "module_path": module_path, 

51 "class_name": class_name, 

52 } 

53 

54 

55class TaskResponse(BaseModel): 

56 """Task serializer for responses.""" 

57 

58 task_id: str | None 

59 task_display_name: str | None 

60 owner: str | None 

61 start_date: datetime | None 

62 end_date: datetime | None 

63 trigger_rule: str | None 

64 depends_on_past: bool 

65 wait_for_downstream: bool 

66 retries: float | None 

67 queue: str | None 

68 pool: str | None 

69 pool_slots: float | None 

70 execution_timeout: TimeDeltaWithValidation | None 

71 retry_delay: TimeDeltaWithValidation | None 

72 retry_exponential_backoff: float 

73 priority_weight: float | None 

74 weight_rule: str | None 

75 ui_color: str | None 

76 ui_fgcolor: str | None 

77 template_fields: list[str] | None 

78 downstream_task_ids: list[str] | None 

79 doc_md: str | None 

80 operator_name: str | None 

81 params: abc.MutableMapping | None 

82 class_ref: dict | None 

83 is_mapped: bool | None 

84 

85 @model_validator(mode="before") 

86 @classmethod 

87 def validate_model(cls, task: Any) -> Any: 

88 task.__dict__.update({"class_ref": _get_class_ref(task), "is_mapped": task.is_mapped}) 

89 return task 

90 

91 @field_validator("weight_rule", mode="before") 

92 @classmethod 

93 def validate_weight_rule(cls, wr: str | PriorityWeightStrategy | None) -> str | None: 

94 """Validate the weight_rule property.""" 

95 if wr is None: 95 ↛ 96line 95 didn't jump to line 96 because the condition on line 95 was never true

96 return None 

97 if isinstance(wr, str): 97 ↛ 98line 97 didn't jump to line 98 because the condition on line 97 was never true

98 return wr 

99 strat_type = type(validate_and_load_priority_weight_strategy(wr)) 

100 try: 

101 return get_weight_rule_from_priority_weight_strategy(strat_type) 

102 except KeyError: 

103 return qualname(strat_type) 

104 

105 @field_validator("params", mode="before") 

106 @classmethod 

107 def get_params(cls, params: SerializedParamsDict | None) -> dict | None: 

108 """Convert params attribute to dict representation.""" 

109 if params is None: 109 ↛ 110line 109 didn't jump to line 110 because the condition on line 109 was never true

110 return None 

111 return {k: v.dump() for k, v in params.items()} 

112 

113 # Mypy issue https://github.com/python/mypy/issues/1362 

114 @computed_field # type: ignore[prop-decorator] 

115 @property 

116 def extra_links(self) -> list[str]: 

117 """Extract and return extra_links.""" 

118 return getattr(self, "operator_extra_links", []) 

119 

120 

121class TaskCollectionResponse(BaseModel): 

122 """Task collection serializer for responses.""" 

123 

124 tasks: list[TaskResponse] 

125 total_entries: int