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
« 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
20from fastapi import Depends, HTTPException, status
21from sqlalchemy.sql import select
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
32extra_links_router = AirflowRouter(
33 tags=["Extra Links"], prefix="/dags/{dag_id}/dagRuns/{dag_run_id}/taskInstances/{task_id}/links"
34)
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
56 dag_run = session.scalar(select(DagRun).where(DagRun.dag_id == dag_id, DagRun.run_id == dag_run_id))
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 )
67 if not ti:
68 raise HTTPException(
69 status.HTTP_404_NOT_FOUND,
70 "TaskInstance not found",
71 )
73 dag = get_dag_for_run_or_latest_version(dag_bag, dag_run, dag_id, session)
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")
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
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)}
108 return ExtraLinkCollectionResponse(
109 extra_links=all_extra_links,
110 total_entries=len(all_extra_links),
111 )