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
« 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
20import inspect
21from collections import abc
22from datetime import datetime
23from typing import TYPE_CHECKING, Any
25from pydantic import computed_field, field_validator, model_validator
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)
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
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
47 class_name = obj.task_type
49 return {
50 "module_path": module_path,
51 "class_name": class_name,
52 }
55class TaskResponse(BaseModel):
56 """Task serializer for responses."""
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
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
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)
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()}
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", [])
121class TaskCollectionResponse(BaseModel):
122 """Task collection serializer for responses."""
124 tasks: list[TaskResponse]
125 total_entries: int