Coverage for /usr/local/lib/python3.10/site-packages/opal_common-0.0.0-py3.10.egg/opal_common/fetcher/engine/fetching_engine.py: 31%

83 statements  

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

1import asyncio 

2import uuid 

3from typing import Coroutine, Dict, List, Union 

4 

5from opal_common.fetcher.engine.base_fetching_engine import BaseFetchingEngine 

6from opal_common.fetcher.engine.core_callbacks import OnFetchFailureCallback 

7from opal_common.fetcher.engine.fetch_worker import fetch_worker 

8from opal_common.fetcher.events import FetcherConfig, FetchEvent 

9from opal_common.fetcher.fetch_provider import BaseFetchProvider 

10from opal_common.fetcher.fetcher_register import FetcherRegister 

11from opal_common.fetcher.logger import get_logger 

12 

13logger = get_logger("engine") 

14 

15 

16class FetchingEngine(BaseFetchingEngine): 

17 """A Task queue manager for fetching events. 

18 

19 - Configure with different fetcher providers - via __init__'s register_config or via self.register.register_fetcher() 

20 - Use queue_url() to fetch a given URL with the default FetchProvider 

21 - Use queue_fetch_event() to fetch data using a configured FetchProvider 

22 - Use with 'async with' to terminate tasks (or call self.terminate_tasks() when done) 

23 """ 

24 

25 DEFAULT_WORKER_COUNT = 6 

26 DEFAULT_CALLBACK_TIMEOUT = 10 

27 DEFAULT_ENQUEUE_TIMEOUT = 10 

28 

29 @staticmethod 

30 def gen_uid(): 

31 return uuid.uuid4().hex 

32 

33 def __init__( 

34 self, 

35 register_config: Dict[str, BaseFetchProvider] = None, 

36 worker_count: int = DEFAULT_WORKER_COUNT, 

37 callback_timeout: int = DEFAULT_CALLBACK_TIMEOUT, 

38 enqueue_timeout: int = DEFAULT_ENQUEUE_TIMEOUT, 

39 retry_config=None, 

40 ) -> None: 

41 # The internal task queue (created at start_workers) 

42 self._queue: asyncio.Queue = None 

43 # Worker working the queue 

44 self._tasks: List[asyncio.Task] = [] 

45 # register of the fetch providers workers can use 

46 self._fetcher_register = FetcherRegister(register_config) 

47 # core event callback registers 

48 self._failure_handlers: List[OnFetchFailureCallback] = [] 

49 # how many workers to run 

50 self._worker_count: int = worker_count 

51 # time in seconds before timeout on a fetch callback 

52 self._callback_timeout = callback_timeout 

53 # time in seconds before time out on adding a task to queue (when full) 

54 self._enqueue_timeout = enqueue_timeout 

55 self._retry_config = retry_config 

56 

57 def start_workers(self): 

58 if self._queue is None: 

59 self._queue = asyncio.Queue() 

60 # create worker tasks 

61 for _ in range(self._worker_count): 

62 self.create_worker() 

63 

64 @property 

65 def register(self) -> FetcherRegister: 

66 return self._fetcher_register 

67 

68 async def __aenter__(self): 

69 """Async Context manager to cancel tasks on exit.""" 

70 self.start_workers() 

71 return self 

72 

73 async def __aexit__(self, exc_type, exc, tb): 

74 if exc is not None: 

75 logger.error( 

76 "Error occurred within FetchingEngine context", 

77 exc_info=repr((exc_type, exc, tb)), 

78 ) 

79 await self.terminate_workers() 

80 

81 async def terminate_workers(self): 

82 """Cancel and wait on the internal worker tasks.""" 

83 # Cancel our worker tasks. 

84 for task in self._tasks: 

85 task.cancel() 

86 # Wait until all worker tasks are cancelled. 

87 await asyncio.gather(*self._tasks, return_exceptions=True) 

88 # reset queue 

89 self._queue = None 

90 

91 async def handle_url(self, url: str, timeout: float = None, **kwargs): 

92 """ 

93 Same as self.queue_url but instead of using a callback, you can wait on this coroutine for the result as a return value 

94 Args: 

95 url (str): 

96 timeout (float, optional): time in seconds to wait on the queued fetch task. Defaults to self._callback_timeout. 

97 kwargs: additional args passed to self.queue_url 

98 

99 Raises: 

100 asyncio.TimeoutError: if the given timeout has expired 

101 also - @see self.queue_fetch_event 

102 """ 

103 timeout = self._callback_timeout if timeout is None else timeout 

104 wait_event = asyncio.Event() 

105 data = {"result": None} 

106 # Callback to wait and retrieve data 

107 

108 async def waiter_callback(answer): 

109 data["result"] = answer 

110 # Signal callback is done 

111 wait_event.set() 

112 

113 await self.queue_url(url, waiter_callback, **kwargs) 

114 # Wait with timeout 

115 if timeout is not None: 

116 await asyncio.wait_for(wait_event.wait(), timeout) 

117 # wait forever 

118 else: 

119 await wait_event.wait() 

120 # return saved result value from callback 

121 return data["result"] 

122 

123 async def queue_url( 

124 self, 

125 url: str, 

126 callback: Coroutine, 

127 config: Union[FetcherConfig, dict, None] = None, 

128 fetcher="HttpFetchProvider", 

129 ) -> FetchEvent: 

130 """Simplified default fetching handler for queuing a fetch task. 

131 

132 Args: 

133 url (str): the URL to fetch from 

134 callback (Coroutine): a callback to call with the fetched result 

135 config (FetcherConfig, optional): Configuration to be used by the fetcher. Defaults to None. 

136 fetcher (str, optional): Which fetcher class to use. Defaults to "HttpFetchProvider". 

137 Returns: 

138 the queued event (which will be mutated to at least have an Id) 

139 

140 Raises: 

141 @see self.queue_fetch_event 

142 """ 

143 # override default fetcher with (potential) override value from FetcherConfig 

144 if isinstance(config, dict) and config.get("fetcher", None) is not None: 

145 fetcher = config["fetcher"] 

146 elif isinstance(config, FetcherConfig) and config.fetcher is not None: 

147 fetcher = config.fetcher 

148 

149 # init a URL event 

150 event = FetchEvent( 

151 url=url, fetcher=fetcher, config=config, retry=self._retry_config 

152 ) 

153 return await self.queue_fetch_event(event, callback) 

154 

155 async def queue_fetch_event( 

156 self, event: FetchEvent, callback: Coroutine, enqueue_timeout=None 

157 ) -> FetchEvent: 

158 """Basic handler to queue a fetch event for a fetcher class. Waits if 

159 the queue is full until enqueue_timeout seconds; if enqueue_timeout is 

160 None returns immediately or raises QueueFull. 

161 

162 Args: 

163 event (FetchEvent): the fetch event to queue as a task 

164 callback (Coroutine): a callback to call with the fetched result 

165 enqueue_timeout (float): timeout in seconds or None for no timeout, Defaults to self.DEFAULT_ENQUEUE_TIMEOUT 

166 

167 Returns: 

168 the queued event (which will be mutated to at least have an Id) 

169 

170 Raises: 

171 asyncio.QueueFull: if the queue is full and enqueue_timeout is set as None 

172 asyncio.TimeoutError: if enqueue_timeout is not None, and the queue is full and hasn't cleared by the timeout time 

173 """ 

174 enqueue_timeout = ( 

175 enqueue_timeout if enqueue_timeout is not None else self._enqueue_timeout 

176 ) 

177 # Assign a unique identifier for the event 

178 event.id = self.gen_uid() 

179 # add to the queue for handling 

180 # if no timeout we return immediately or raise QueueFull 

181 if enqueue_timeout is None: 

182 await self._queue.put_nowait((event, callback)) 

183 # if timeout 

184 else: 

185 await asyncio.wait_for(self._queue.put((event, callback)), enqueue_timeout) 

186 return event 

187 

188 def create_worker(self) -> asyncio.Task: 

189 """Create an asyncio worker to work the engine's queue Engine init 

190 starts several workers according to given configuration.""" 

191 task = asyncio.create_task(fetch_worker(self._queue, self)) 

192 self._tasks.append(task) 

193 return task 

194 

195 def register_failure_handler(self, callback: OnFetchFailureCallback): 

196 """Register a callback to be called with exception and original event 

197 in case of failure. 

198 

199 Args: 

200 callback (OnFetchFailureCallback): callback to register 

201 """ 

202 self._failure_handlers.append(callback) 

203 

204 async def _management_event_handler( 

205 self, handlers: List[Coroutine], *args, **kwargs 

206 ): 

207 """ 

208 Generic management event subscriber caller 

209 Args: 

210 handlers (List[Coroutine]): callback coroutines 

211 """ 

212 await asyncio.gather(*(callback(*args, **kwargs) for callback in handlers)) 

213 

214 async def _on_failure(self, error: Exception, event: FetchEvent): 

215 """Call event failure subscribers. 

216 

217 Args: 

218 error (Exception): thrown exception 

219 event (FetchEvent): event which was being handled 

220 """ 

221 await self._management_event_handler(self._failure_handlers, error, event)