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
« prev ^ index » next coverage.py v7.15.2, created at 2026-10-10 11:54 +0000
1import asyncio
2from typing import Coroutine
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
9logger = get_logger("fetch_worker")
12async def fetch_worker(queue: asyncio.Queue, engine):
13 """The worker task performing items added to the Engine's Queue.
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()