Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/common/partition_helpers.py: 31%

29 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. 

17from __future__ import annotations 

18 

19from typing import TYPE_CHECKING 

20 

21import structlog 

22 

23from airflow.exceptions import DeserializationError 

24from airflow.models.serialized_dag import SerializedDagModel 

25from airflow.timetables.simple import PartitionedAssetTimetable 

26 

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

28 from sqlalchemy.orm import Session 

29 

30 

31log = structlog.get_logger(logger_name=__name__) 

32 

33 

34def _extract_partitioned_timetable(serdag: SerializedDagModel) -> PartitionedAssetTimetable | None: 

35 """Return the ``PartitionedAssetTimetable`` carried by *serdag*, or ``None``.""" 

36 try: 

37 timetable = serdag.dag.timetable 

38 except (DeserializationError, TypeError): 

39 # ``DeserializationError`` covers structural serialization failures so 

40 # a corrupted serialized Dag silently degrades to non-partitioned 

41 # rather than 500-ing the read-only UI page. ``TypeError`` covers the 

42 # eager validation in ``RollupMapper.__init__`` (raised during 

43 # timetable deserialization when an upstream mapper / window pair is 

44 # incompatible), so a misconfigured rollup Dag also degrades to 

45 # non-partitioned here rather than 500-ing. 

46 # ``KeyError`` / ``AttributeError`` / ``ImportError`` are intentionally 

47 # not caught: refactor bugs that rename an attribute on ``serdag.dag`` 

48 # or break an import path must surface to the caller rather than 

49 # silently downgrading the route to non-rollup. 

50 log.warning("Failed to deserialize timetable for Dag", dag_id=serdag.dag_id, exc_info=True) 

51 return None 

52 if not timetable.partitioned: 

53 return None 

54 if TYPE_CHECKING: 

55 assert isinstance(timetable, PartitionedAssetTimetable) 

56 return timetable 

57 

58 

59def load_partitioned_timetable(dag_id: str, session: Session) -> PartitionedAssetTimetable | None: 

60 """ 

61 Return the PartitionedAssetTimetable for *dag_id*, or None if absent or not partitioned. 

62 

63 Callers gate this behind ``DagModel.has_rollup_mappers``, which is only 

64 populated for ``PartitionedAssetTimetable``. The ``TYPE_CHECKING`` assert 

65 narrows the type for mypy without a runtime ``isinstance`` cost. 

66 """ 

67 serdag = SerializedDagModel.get(dag_id=dag_id, session=session) 

68 if serdag is None: 

69 return None 

70 return _extract_partitioned_timetable(serdag) 

71 

72 

73def load_partitioned_timetables( 

74 dag_ids: list[str], session: Session 

75) -> dict[str, PartitionedAssetTimetable | None]: 

76 """ 

77 Batch-load PartitionedAssetTimetables for *dag_ids* in a single query. 

78 

79 Routes that already gate per-Dag on ``DagModel.has_rollup_mappers`` should 

80 use this when iterating over many Dags so ``SerializedDagModel`` is hit 

81 once instead of once per Dag. Returns a dict keyed by ``dag_id``; entries 

82 whose timetable failed to deserialize or is not partitioned are ``None``. 

83 """ 

84 if not dag_ids: 

85 return {} 

86 return { 

87 serdag.dag_id: _extract_partitioned_timetable(serdag) 

88 for serdag in SerializedDagModel.get_latest_serialized_dags(dag_ids=dag_ids, session=session) 

89 }