Coverage for /usr/local/lib/python3.10/site-packages/opal_common-0.0.0-py3.10.egg/opal_common/fetcher/providers/http_fetch_provider.py: 46%

85 statements  

« prev     ^ index     » next       coverage.py v7.15.2, created at 2026-10-10 11:54 +0000

1"""Simple HTTP get data fetcher using requests supports.""" 

2 

3from enum import Enum 

4from typing import Any, ClassVar, Set, Union, cast 

5 

6import httpx 

7from aiohttp import ClientResponse, ClientSession, ClientTimeout 

8from opal_common.config import opal_common_config 

9from opal_common.fetcher.events import FetcherConfig, FetchEvent 

10from opal_common.fetcher.fetch_provider import BaseFetchProvider 

11from opal_common.fetcher.logger import get_logger 

12from opal_common.http_utils import is_http_error_response, redact_url 

13from opal_common.security.sslcontext import get_custom_ssl_context 

14from pydantic import validator 

15 

16logger = get_logger("http_fetch_provider") 

17 

18 

19class HttpMethods(Enum): 

20 GET = "get" 

21 POST = "post" 

22 PUT = "put" 

23 PATCH = "patch" 

24 HEAD = "head" 

25 DELETE = "delete" 

26 

27 

28class HttpFetcherConfig(FetcherConfig): 

29 """Config for HttpFetchProvider's Adding HTTP headers.""" 

30 

31 # ``headers`` carries Authorization tokens and ``data`` the (possibly 

32 # sensitive) payload - mask both in repr/str so they never leak into logs. 

33 _redacted_repr_fields: ClassVar[Set[str]] = {"headers", "data"} 

34 

35 headers: dict = None 

36 is_json: bool = True 

37 process_data: bool = True 

38 method: HttpMethods = HttpMethods.GET 

39 data: Any = None 

40 

41 @validator("method") 

42 def force_enum(cls, v): 

43 if isinstance(v, str): 43 ↛ 45line 43 didn't jump to line 45 because the condition on line 43 was always true

44 return HttpMethods(v) 

45 if isinstance(v, HttpMethods): 

46 return v 

47 raise ValueError(f"invalid value: {v}") 

48 

49 class Config: 

50 use_enum_values = True 

51 

52 

53class HttpFetchEvent(FetchEvent): 

54 fetcher: str = "HttpFetchProvider" 

55 config: HttpFetcherConfig = None 

56 

57 

58class HttpFetchProvider(BaseFetchProvider): 

59 def __init__(self, event: HttpFetchEvent) -> None: 

60 self._event: HttpFetchEvent 

61 if event.config is None: 

62 event.config = HttpFetcherConfig() 

63 super().__init__(event) 

64 self._session = None 

65 self._custom_ssl_context = get_custom_ssl_context() 

66 self._ssl_context_kwargs = ( 

67 {"ssl": self._custom_ssl_context} 

68 if self._custom_ssl_context is not None 

69 else {} 

70 ) 

71 

72 def parse_event(self, event: FetchEvent) -> HttpFetchEvent: 

73 return HttpFetchEvent(**event.dict(exclude={"config"}), config=event.config) 

74 

75 async def __aenter__(self): 

76 headers = {} 

77 timeout = opal_common_config.HTTP_FETCHER_TIMEOUT 

78 if self._event.config.headers is not None: 

79 headers = self._event.config.headers 

80 if opal_common_config.HTTP_FETCHER_PROVIDER_CLIENT == "httpx": 

81 self._session = httpx.AsyncClient( 

82 headers=headers, timeout=timeout, trust_env=True 

83 ) 

84 else: 

85 self._session = ClientSession( 

86 headers=headers, 

87 raise_for_status=True, 

88 timeout=ClientTimeout(total=timeout), 

89 trust_env=True, 

90 ) 

91 self._session = await self._session.__aenter__() 

92 return self 

93 

94 async def __aexit__(self, exc_type=None, exc_val=None, tb=None): 

95 await self._session.__aexit__(exc_type, exc_val, tb) 

96 

97 async def _fetch_(self): 

98 logger.debug(f"{self.__class__.__name__} fetching from {redact_url(self._url)}") 

99 http_method = self.match_http_method_from_type( 

100 self._session, self._event.config.method 

101 ) 

102 if self._event.config.data is not None: 

103 result: Union[ClientResponse, httpx.Response] = await http_method( 

104 self._url, data=self._event.config.data, **self._ssl_context_kwargs 

105 ) 

106 else: 

107 result = await http_method(self._url, **self._ssl_context_kwargs) 

108 result.raise_for_status() 

109 return result 

110 

111 @staticmethod 

112 def match_http_method_from_type( 

113 session: Union[ClientSession, httpx.AsyncClient], method_type: HttpMethods 

114 ): 

115 return getattr(session, method_type.value) 

116 

117 @staticmethod 

118 async def _response_to_data( 

119 res: Union[ClientResponse, httpx.Response], *, is_json: bool 

120 ) -> Any: 

121 if isinstance(res, httpx.Response): 

122 return res.json() if is_json else res.text 

123 else: 

124 res = cast(ClientResponse, res) 

125 return await (res.json() if is_json else res.text()) 

126 

127 async def _process_(self, res: Union[ClientResponse, httpx.Response]): 

128 # do not process data when the http response is an error 

129 if is_http_error_response(res): 

130 return res 

131 

132 # if we are asked to process the data before we return it 

133 if self._event.config.process_data: 

134 data = await self._response_to_data(res, is_json=self._event.config.is_json) 

135 return data 

136 # return raw result 

137 else: 

138 return res