Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/core_api/routes/public/job.py: 100%
21 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.
17from __future__ import annotations
19from typing import Annotated
21from fastapi import Depends, status
22from sqlalchemy import select
23from sqlalchemy.orm import joinedload
25from airflow.api_fastapi.common.db.common import (
26 SessionDep,
27 paginated_select,
28)
29from airflow.api_fastapi.common.parameters import (
30 FilterParam,
31 QueryLimit,
32 QueryOffset,
33 RangeFilter,
34 SortParam,
35 datetime_range_filter_factory,
36 filter_param_factory,
37)
38from airflow.api_fastapi.common.router import AirflowRouter
39from airflow.api_fastapi.core_api.datamodels.job import (
40 JobCollectionResponse,
41)
42from airflow.api_fastapi.core_api.openapi.exceptions import create_openapi_http_exception_doc
43from airflow.api_fastapi.core_api.security import AccessView, requires_access_view
44from airflow.jobs.job import Job
46job_router = AirflowRouter(tags=["Job"], prefix="/jobs")
49@job_router.get(
50 "",
51 responses=create_openapi_http_exception_doc([status.HTTP_400_BAD_REQUEST]),
52 dependencies=[Depends(requires_access_view(AccessView.JOBS))],
53)
54def get_jobs(
55 start_date_range: Annotated[
56 RangeFilter,
57 Depends(datetime_range_filter_factory("start_date", Job)),
58 ],
59 end_date_range: Annotated[
60 RangeFilter,
61 Depends(datetime_range_filter_factory("end_date", Job)),
62 ],
63 limit: QueryLimit,
64 offset: QueryOffset,
65 order_by: Annotated[
66 SortParam,
67 Depends(
68 SortParam(
69 [
70 "id",
71 "dag_id",
72 "state",
73 "job_type",
74 "start_date",
75 "end_date",
76 "latest_heartbeat",
77 "executor_class",
78 "hostname",
79 "unixname",
80 ],
81 Job,
82 ).dynamic_depends(default="id")
83 ),
84 ],
85 session: SessionDep,
86 state: Annotated[
87 FilterParam[str | None], Depends(filter_param_factory(Job.state, str | None, filter_name="job_state"))
88 ],
89 job_type: Annotated[
90 FilterParam[str | None],
91 Depends(filter_param_factory(Job.job_type, str | None, filter_name="job_type")),
92 ],
93 hostname: Annotated[
94 FilterParam[str | None],
95 Depends(filter_param_factory(Job.hostname, str | None, filter_name="hostname")),
96 ],
97 executor_class: Annotated[
98 FilterParam[str | None],
99 Depends(filter_param_factory(Job.executor_class, str | None, filter_name="executor_class")),
100 ],
101 is_alive: bool | None = None,
102) -> JobCollectionResponse:
103 """Get all jobs."""
104 base_select = select(Job).order_by(Job.latest_heartbeat.desc()).options(joinedload(Job.dag_model))
106 jobs_select, total_entries = paginated_select(
107 statement=base_select,
108 filters=[
109 start_date_range,
110 end_date_range,
111 state,
112 job_type,
113 hostname,
114 executor_class,
115 ],
116 order_by=order_by,
117 limit=limit,
118 offset=offset,
119 session=session,
120 return_total_entries=True,
121 )
122 jobs = session.scalars(jobs_select).all()
124 if is_alive is not None:
125 jobs = [job for job in jobs if job.is_alive()]
127 return JobCollectionResponse(
128 jobs=jobs,
129 total_entries=total_entries,
130 )