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

25 statements  

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

1import asyncio 

2from typing import Coroutine 

3 

4from opal_common.fetcher.engine.base_fetching_engine import BaseFetchingEngine 

5from opal_common.fetcher.events import FetchEvent 

6from opal_common.fetcher.fetcher_register import FetcherRegister 

7from opal_common.fetcher.logger import get_logger 

8 

9logger = get_logger("fetch_worker") 

10 

11 

12async def fetch_worker(queue: asyncio.Queue, engine): 

13 """The worker task performing items added to the Engine's Queue. 

14 

15 Args: 

16 queue (asyncio.Queue): The Queue 

17 engine (BaseFetchingEngine): The engine itself 

18 """ 

19 engine: BaseFetchingEngine 

20 register: FetcherRegister = engine.register 

21 while True: 

22 # types 

23 event: FetchEvent 

24 callback: Coroutine 

25 # get a event from the queue 

26 event, callback = await queue.get() 

27 # take care of it 

28 try: 

29 # get fetcher for the event 

30 fetcher = register.get_fetcher_for_event(event) 

31 # fetch 

32 async with fetcher: 

33 res = await fetcher.fetch() 

34 data = await fetcher.process(res) 

35 # callback to event owner 

36 try: 

37 await callback(data) 

38 except Exception as err: 

39 logger.exception(f"Fetcher callback - {callback} failed") 

40 await engine._on_failure(err, event) 

41 except Exception as err: 

42 logger.exception("Failed to process fetch event") 

43 await engine._on_failure(err, event) 

44 finally: 

45 # Notify the queue that the "work item" has been processed. 

46 queue.task_done()