Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/extra_links.py: 78%

34 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 fastapi import Depends, HTTPException, status 

21from sqlalchemy.sql import select 

22 

23from airflow.api_fastapi.common.dagbag import DagBagDep, get_dag_for_run_or_latest_version 

24from airflow.api_fastapi.common.db.common import SessionDep 

25from airflow.api_fastapi.common.router import AirflowRouter 

26from airflow.api_fastapi.core_api.datamodels.extra_links import ExtraLinkCollectionResponse 

27from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc 

28from airflow.api_fastapi.core_api.security import DagAccessEntity, requires_access_dag 

29from airflow.exceptions import TaskNotFound 

30from airflow.models import DagRun 

31 

32extra_links_router = AirflowRouter( 

33 tags=["Extra Links"], prefix="/dags/{dag_id}/dagRuns/{dag_run_id}/taskInstances/{task_id}/links" 

34) 

35 

36 

37@extra_links_router.get( 

38 "", 

39 responses=create_openapi_http_exception_doc([status.HTTP_404_NOT_FOUND]), 

40 dependencies=[Depends(requires_access_dag("GET", DagAccessEntity.TASK_INSTANCE))], 

41 tags=["Task Instance"], 

42) 

43def get_extra_links( 

44 dag_id: str, 

45 dag_run_id: str, 

46 task_id: str, 

47 session: SessionDep, 

48 dag_bag: DagBagDep, 

49 map_index: int = -1, 

50 try_number: int | None = None, 

51) -> ExtraLinkCollectionResponse: 

52 """Get extra links for task instance.""" 

53 from airflow.models.taskinstance import TaskInstance 

54 from airflow.models.taskinstancehistory import TaskInstanceHistory 

55 

56 dag_run = session.scalar(select(DagRun).where(DagRun.dag_id == dag_id, DagRun.run_id == dag_run_id)) 

57 

58 ti = session.scalar( 

59 select(TaskInstance).where( 

60 TaskInstance.dag_id == dag_id, 

61 TaskInstance.run_id == dag_run_id, 

62 TaskInstance.task_id == task_id, 

63 TaskInstance.map_index == map_index, 

64 ) 

65 ) 

66 

67 if not ti: 

68 raise HTTPException( 

69 status.HTTP_404_NOT_FOUND, 

70 "TaskInstance not found", 

71 ) 

72 

73 dag = get_dag_for_run_or_latest_version(dag_bag, dag_run, dag_id, session) 

74 

75 try: 

76 task = dag.get_task(task_id) 

77 except TaskNotFound: 

78 raise HTTPException(status.HTTP_404_NOT_FOUND, f"Task with ID = {task_id} not found") 

79 

80 # Resolve which object to use for link generation. For the current try we use 

81 # the live TI; for past tries we fetch the immutable TaskInstanceHistory record, 

82 # which also validates that the requested try_number actually exists. 

83 if try_number is not None and try_number != ti.try_number: 83 ↛ 84line 83 didn't jump to line 84 because the condition on line 83 was never true

84 tih = session.scalar( 

85 select(TaskInstanceHistory).where( 

86 TaskInstanceHistory.dag_id == dag_id, 

87 TaskInstanceHistory.task_id == task_id, 

88 TaskInstanceHistory.run_id == dag_run_id, 

89 TaskInstanceHistory.map_index == map_index, 

90 TaskInstanceHistory.try_number == try_number, 

91 ) 

92 ) 

93 if not tih: 

94 raise HTTPException( 

95 status.HTTP_404_NOT_FOUND, 

96 f"TaskInstanceHistory not found for try_number={try_number}", 

97 ) 

98 ti_for_links = tih 

99 else: 

100 ti_for_links = ti 

101 

102 all_extra_link_pairs = ( 

103 (link_name, task.get_extra_links(ti_for_links, link_name)) 

104 for link_name in task.extra_links # type: ignore[arg-type] 

105 ) 

106 all_extra_links = {link_name: link_url or None for link_name, link_url in sorted(all_extra_link_pairs)} 

107 

108 return ExtraLinkCollectionResponse( 

109 extra_links=all_extra_links, 

110 total_entries=len(all_extra_links), 

111 )