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

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 uuid import UUID 

20 

21from cadwyn import VersionedAPIRouter 

22from fastapi import HTTPException, Security, status 

23 

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) 

36 

37router = VersionedAPIRouter( 

38 route_class=ExecutionAPIRoute, 

39 dependencies=[ 

40 Security(require_auth, scopes=["ct:self", "token:workload"]), 

41 ], 

42) 

43 

44 

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 ) 

68 

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 ) 

80 

81 ct.state = ConnectionTestState.RUNNING 

82 

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 ) 

93 

94 

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 ) 

118 

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 ) 

135 

136 ct.state = body.state 

137 ct.result_message = body.result_message 

138 

139 if body.state == ConnectionTestState.SUCCESS and ct.commit_on_success: 

140 ct.commit_to_connection_table(session=session)