Coverage for /home/airflow/.local/lib/python3.12/site-packages/airflow/api_fastapi/common/http_access_log.py: 87%

59 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"""HTTP access log middleware using structlog.""" 

18 

19from __future__ import annotations 

20 

21import contextlib 

22import time 

23from typing import TYPE_CHECKING 

24from urllib.parse import parse_qsl, urlencode 

25 

26import structlog 

27 

28from airflow._shared.secrets_masker import secrets_masker 

29 

30if TYPE_CHECKING: 30 ↛ 31line 30 didn't jump to line 31 because the condition on line 30 was never true

31 from starlette.types import ASGIApp, Message, Receive, Scope, Send 

32 

33logger = structlog.get_logger(logger_name="http.access") 

34 

35_HEALTH_PATHS = frozenset(["/api/v2/monitor/health"]) 

36 

37 

38def _redact_query_string(query: str) -> str: 

39 """ 

40 Redact secret-looking query parameters before they reach the access log. 

41 

42 Treat each ``key=value`` pair independently so a key whose name signals a secret 

43 (``password``, ``token``, ``api_key`` — anything ``secrets_masker`` flags as sensitive) 

44 gets its value replaced with ``***``. Also catches values that were previously registered 

45 via ``mask_secret()``. 

46 """ 

47 if not query: 

48 return query 

49 try: 

50 pairs = parse_qsl(query, keep_blank_values=True) 

51 except ValueError: 

52 # Malformed query string — leave it alone; we'd rather log the raw bytes than 

53 # silently drop diagnostic information. 

54 return query 

55 if not pairs: 55 ↛ 56line 55 didn't jump to line 56 because the condition on line 55 was never true

56 return query 

57 redacted_pairs = [(k, secrets_masker.redact(v, k)) for k, v in pairs] 

58 return urlencode(redacted_pairs) 

59 

60 

61class HttpAccessLogMiddleware: 

62 """ 

63 Log completed HTTP requests as structured log events. 

64 

65 This middleware replaces uvicorn's built-in access logger. It measures the 

66 full round-trip duration, binds any ``x-request-id`` header value to the 

67 structlog context for the duration of the request, and emits one log event 

68 per completed request. 

69 

70 Health-check paths are excluded to avoid log noise. 

71 """ 

72 

73 def __init__( 

74 self, 

75 app: ASGIApp, 

76 request_id_header: str = "x-request-id", 

77 ) -> None: 

78 self.app = app 

79 self.request_id_header = request_id_header.lower().encode("ascii") 

80 

81 async def __call__(self, scope: Scope, receive: Receive, send: Send) -> None: 

82 if scope["type"] != "http": 

83 await self.app(scope, receive, send) 

84 return 

85 

86 start = time.monotonic_ns() 

87 response: Message | None = None 

88 

89 async def capture_send(message: Message) -> None: 

90 nonlocal response 

91 if message["type"] == "http.response.start": 

92 response = message 

93 await send(message) 

94 

95 request_id: str | None = None 

96 for name, value in scope["headers"]: 

97 if name == self.request_id_header: 97 ↛ 98line 97 didn't jump to line 98 because the condition on line 97 was never true

98 request_id = value.decode("ascii", errors="replace") 

99 break 

100 

101 ctx = ( 

102 structlog.contextvars.bound_contextvars(request_id=request_id) 

103 if request_id is not None 

104 else contextlib.nullcontext() 

105 ) 

106 

107 with ctx: 

108 try: 

109 await self.app(scope, receive, capture_send) 

110 except Exception: 

111 if response is None: 

112 response = {"status": 500} 

113 raise 

114 finally: 

115 path = scope["path"] 

116 if path not in _HEALTH_PATHS: 

117 duration_us = (time.monotonic_ns() - start) // 1000 

118 status = response["status"] if response is not None else 0 

119 method = scope.get("method", "") 

120 query = _redact_query_string(scope["query_string"].decode("ascii", errors="replace")) 

121 client = scope.get("client") 

122 client_addr = f"{client[0]}:{client[1]}" if client else None 

123 

124 # Guard the log emit: if it raised inside a ``finally`` while the 

125 # original ``try`` block was already propagating an app exception, 

126 # Python's exception-replacement semantics would discard the 

127 # original. Swallow logging failures so the application exception 

128 # always reaches uvicorn intact. 

129 with contextlib.suppress(Exception): 

130 logger.info( 

131 "request finished", 

132 method=method, 

133 path=path, 

134 query=query, 

135 status_code=status, 

136 duration_us=duration_us, 

137 client_addr=client_addr, 

138 )