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

1from opal_common.fetcher.events import FetchEvent 

2from opal_common.fetcher.logger import get_logger 

3from tenacity import retry, stop, wait 

4 

5logger = get_logger("opal.providers") 

6 

7 

8class BaseFetchProvider: 

9 """Base class for data fetching providers. 

10 

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 """ 

15 

16 DEFAULT_RETRY_CONFIG = { 

17 "wait": wait.wait_random_exponential(), 

18 "stop": stop.stop_after_attempt(200), 

19 "reraise": True, 

20 } 

21 

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 ) 

34 

35 def parse_event(self, event: FetchEvent) -> FetchEvent: 

36 """Parse the event (And config within it) into the right object type. 

37 

38 Args: 

39 event (FetchEvent): the event to be parsed 

40 

41 Returns: 

42 FetchEvent: an event deriving from FetchEvent 

43 """ 

44 return event 

45 

46 async def fetch(self): 

47 """Fetch and return data. 

48 

49 Calls self._fetch_ with a retry mechanism 

50 """ 

51 attempter = retry(**self._retry_config)(self._fetch_) 

52 res = await attempter() 

53 return res 

54 

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 

61 

62 async def __aenter__(self): 

63 return self 

64 

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

66 pass 

67 

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 

72 

73 async def _process_(self, data): 

74 return data 

75 

76 def set_retry_config(self, retry_config: dict): 

77 """Set the configuration for retrying failed fetches. 

78 

79 @see self.DEFAULT_RETRY_CONFIG 

80 

81 Args: 

82 retry_config (dict): Tenacity retry config 

83 """ 

84 self._retry_config = retry_config