morington

Untitled

Jan 24th, 2024
614
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
Python 1.77 KB | None | 0 0
  1. import asyncio
  2. from concurrent.futures import ThreadPoolExecutor
  3. from time import sleep
  4. from typing import Callable, List
  5.  
  6. import structlog
  7.  
  8. import setup_logger
  9.  
  10. logger = structlog.get_logger(__name__)
  11.  
  12.  
  13. class Emitter:
  14.     def __init__(self):
  15.         self.handlers: List[Callable] = []
  16.  
  17.     def connect(self, handler: Callable):
  18.         self.handlers.append(handler)
  19.  
  20.     def emit(self, *args, **kwargs):
  21.         for handler in self.handlers:
  22.             handler(*args, **kwargs)
  23.  
  24.  
  25. class MySignal:
  26.     def __init__(self):
  27.         self._emitter = Emitter()
  28.  
  29.     def connect(self, handler: Callable):
  30.         self._emitter.connect(handler)
  31.  
  32.     def emit(self, *args, **kwargs):
  33.         self._emitter.emit(*args, **kwargs)
  34.  
  35.  
  36. class Handler:
  37.     def __init__(self):
  38.         self.signal = MySignal()
  39.  
  40.     def worker(self, name: str, timeout: int):
  41.         sleep(timeout)
  42.         self.signal.emit(name)
  43.  
  44.  
  45. async def async_generator():
  46.     # Имитация NATS Subscribe
  47.     for i, data in enumerate([
  48.         ("task1", 3), ("task2", 2), ("task3", 4), ("task4", 3), ("task5", 3),
  49.         ("exit", 10)], start=1
  50.     ):
  51.         if data[0] == "exit":
  52.             await asyncio.sleep(data[1])
  53.             continue
  54.  
  55.         await asyncio.sleep(1)
  56.         yield i, data
  57.  
  58.  
  59. def output_handler(name: str):
  60.     logger.debug("<=== Work completed", name=name)
  61.  
  62.  
  63. async def main():
  64.     with ThreadPoolExecutor(max_workers=3) as executor:
  65.         async for i, data in async_generator():
  66.             logger.debug("===> New task", name=data[0], timeout=data[1])
  67.  
  68.             handler = Handler()
  69.             handler.signal.connect(output_handler)
  70.             executor.submit(handler.worker, data[0], data[1])
  71.  
  72.  
  73. if __name__ == '__main__':
  74.     asyncio.run(main())
  75.  
Advertisement
Add Comment
Please, Sign In to add comment