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
« 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
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
13logger = get_logger("engine")
16class FetchingEngine(BaseFetchingEngine):
17 """A Task queue manager for fetching events.
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 """
25 DEFAULT_WORKER_COUNT = 6
26 DEFAULT_CALLBACK_TIMEOUT = 10
27 DEFAULT_ENQUEUE_TIMEOUT = 10
29 @staticmethod
30 def gen_uid():
31 return uuid.uuid4().hex
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
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()
64 @property
65 def register(self) -> FetcherRegister:
66 return self._fetcher_register
68 async def __aenter__(self):
69 """Async Context manager to cancel tasks on exit."""
70 self.start_workers()
71 return self
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()
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
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
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
108 async def waiter_callback(answer):
109 data["result"] = answer
110 # Signal callback is done
111 wait_event.set()
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"]
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.
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)
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
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)
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.
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
167 Returns:
168 the queued event (which will be mutated to at least have an Id)
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
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
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.
199 Args:
200 callback (OnFetchFailureCallback): callback to register
201 """
202 self._failure_handlers.append(callback)
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))
214 async def _on_failure(self, error: Exception, event: FetchEvent):
215 """Call event failure subscribers.
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)