Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/execution_api/routes/connection_tests.py: 30%
31 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 uuid import UUID
21from cadwyn import VersionedAPIRouter
22from fastapi import HTTPException, Security, status
24from airflow.api_fastapi.common.db.common import SessionDep
25from airflow.api_fastapi.execution_api.datamodels.connection_test import (
26 ConnectionTestConnectionResponse,
27 ConnectionTestResultBody,
28)
29from airflow.api_fastapi.execution_api.security import ExecutionAPIRoute, require_auth
30from airflow.models.connection_test import (
31 ACTIVE_STATES,
32 TERMINAL_STATES,
33 ConnectionTestRequest,
34 ConnectionTestState,
35)
37router = VersionedAPIRouter(
38 route_class=ExecutionAPIRoute,
39 dependencies=[
40 Security(require_auth, scopes=["ct:self", "token:workload"]),
41 ],
42)
45@router.get(
46 "/{connection_test_id}/connection",
47 responses={
48 status.HTTP_404_NOT_FOUND: {"description": "Connection test not found"},
49 status.HTTP_409_CONFLICT: {
50 "description": "Connection test already in RUNNING or terminal state",
51 },
52 },
53)
54def get_connection_test_connection(
55 connection_test_id: UUID,
56 session: SessionDep,
57) -> ConnectionTestConnectionResponse:
58 """Return the test request's connection data and atomically mark it RUNNING (single-fetch)."""
59 ct = session.get(ConnectionTestRequest, connection_test_id, with_for_update=True)
60 if ct is None:
61 raise HTTPException(
62 status_code=status.HTTP_404_NOT_FOUND,
63 detail={
64 "reason": "not_found",
65 "message": f"Connection test {connection_test_id} not found",
66 },
67 )
69 if ct.state not in (ConnectionTestState.PENDING, ConnectionTestState.QUEUED):
70 raise HTTPException(
71 status_code=status.HTTP_409_CONFLICT,
72 detail={
73 "reason": "conflict",
74 "message": (
75 f"Connection test {connection_test_id} is in state {ct.state}; "
76 "credentials can only be fetched once while PENDING or QUEUED."
77 ),
78 },
79 )
81 ct.state = ConnectionTestState.RUNNING
83 return ConnectionTestConnectionResponse(
84 conn_id=ct.connection_id,
85 conn_type=ct.conn_type,
86 host=ct.host,
87 login=ct.login,
88 password=ct.password,
89 schema=ct.schema,
90 port=ct.port,
91 extra=ct.extra,
92 )
95@router.patch(
96 "/{connection_test_id}",
97 status_code=status.HTTP_204_NO_CONTENT,
98 responses={
99 status.HTTP_404_NOT_FOUND: {"description": "Connection test not found"},
100 status.HTTP_409_CONFLICT: {"description": "Connection test already in a terminal state"},
101 },
102)
103def patch_connection_test(
104 connection_test_id: UUID,
105 body: ConnectionTestResultBody,
106 session: SessionDep,
107) -> None:
108 """Update the result of a connection test."""
109 ct = session.get(ConnectionTestRequest, connection_test_id, with_for_update=True)
110 if ct is None:
111 raise HTTPException(
112 status_code=status.HTTP_404_NOT_FOUND,
113 detail={
114 "reason": "not_found",
115 "message": f"Connection test {connection_test_id} not found",
116 },
117 )
119 if ct.state in TERMINAL_STATES:
120 raise HTTPException(
121 status_code=status.HTTP_409_CONFLICT,
122 detail={
123 "reason": "conflict",
124 "message": (f"Connection test {connection_test_id} is already in terminal state: {ct.state}"),
125 },
126 )
127 if ct.state not in ACTIVE_STATES:
128 raise HTTPException(
129 status_code=status.HTTP_409_CONFLICT,
130 detail={
131 "reason": "conflict",
132 "message": f"Connection test {connection_test_id} is not in an active state: {ct.state}",
133 },
134 )
136 ct.state = body.state
137 ct.result_message = body.result_message
139 if body.state == ConnectionTestState.SUCCESS and ct.commit_on_success:
140 ct.commit_to_connection_table(session=session)