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
« 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."""
3from enum import Enum
4from typing import Any, ClassVar, Set, Union, cast
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
16logger = get_logger("http_fetch_provider")
19class HttpMethods(Enum):
20 GET = "get"
21 POST = "post"
22 PUT = "put"
23 PATCH = "patch"
24 HEAD = "head"
25 DELETE = "delete"
28class HttpFetcherConfig(FetcherConfig):
29 """Config for HttpFetchProvider's Adding HTTP headers."""
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"}
35 headers: dict = None
36 is_json: bool = True
37 process_data: bool = True
38 method: HttpMethods = HttpMethods.GET
39 data: Any = None
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}")
49 class Config:
50 use_enum_values = True
53class HttpFetchEvent(FetchEvent):
54 fetcher: str = "HttpFetchProvider"
55 config: HttpFetcherConfig = None
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 )
72 def parse_event(self, event: FetchEvent) -> HttpFetchEvent:
73 return HttpFetchEvent(**event.dict(exclude={"config"}), config=event.config)
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
94 async def __aexit__(self, exc_type=None, exc_val=None, tb=None):
95 await self._session.__aexit__(exc_type, exc_val, tb)
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
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)
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())
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
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