Coverage for /usr/local/lib/python3.10/site-packages/opal_common-0.0.0-py3.10.egg/opal_common/fetcher/fetch_provider.py: 47%
32 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
1from opal_common.fetcher.events import FetchEvent
2from opal_common.fetcher.logger import get_logger
3from tenacity import retry, stop, wait
5logger = get_logger("opal.providers")
8class BaseFetchProvider:
9 """Base class for data fetching providers.
11 - Override self._fetch_ to implement fetching
12 - call self.fetch() to retrieve data (wrapped in retries and safe execution guards)
13 - override __aenter__ and __aexit__ for async context
14 """
16 DEFAULT_RETRY_CONFIG = {
17 "wait": wait.wait_random_exponential(),
18 "stop": stop.stop_after_attempt(200),
19 "reraise": True,
20 }
22 def __init__(self, event: FetchEvent, retry_config=None) -> None:
23 """
24 Args:
25 event (FetchEvent): the event desciring what we should fetch
26 retry_config (dict): Tenacity.retry config (@see https://tenacity.readthedocs.io/en/latest/api.html#retry-main-api) for retrying fetching
27 """
28 # convert the event as needed and save it
29 self._event = self.parse_event(event)
30 self._url = event.url
31 self._retry_config = (
32 retry_config if retry_config is not None else self.DEFAULT_RETRY_CONFIG
33 )
35 def parse_event(self, event: FetchEvent) -> FetchEvent:
36 """Parse the event (And config within it) into the right object type.
38 Args:
39 event (FetchEvent): the event to be parsed
41 Returns:
42 FetchEvent: an event deriving from FetchEvent
43 """
44 return event
46 async def fetch(self):
47 """Fetch and return data.
49 Calls self._fetch_ with a retry mechanism
50 """
51 attempter = retry(**self._retry_config)(self._fetch_)
52 res = await attempter()
53 return res
55 async def process(self, data):
56 try:
57 return await self._process_(data)
58 except:
59 logger.exception("Failed to process fetched data")
60 raise
62 async def __aenter__(self):
63 return self
65 async def __aexit__(self, exc_type=None, exc_val=None, tb=None):
66 pass
68 async def _fetch_(self):
69 """Internal fetch operation called by self.fetch() Override this method
70 to implement a new fetch provider."""
71 pass
73 async def _process_(self, data):
74 return data
76 def set_retry_config(self, retry_config: dict):
77 """Set the configuration for retrying failed fetches.
79 @see self.DEFAULT_RETRY_CONFIG
81 Args:
82 retry_config (dict): Tenacity retry config
83 """
84 self._retry_config = retry_config