morington

Untitled

Jan 24th, 2024 (edited)
435
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
Python 6.14 KB | None | 0 0
  1. import asyncio
  2. import threading
  3. from enum import StrEnum
  4. from time import sleep
  5. from typing import Callable, List, Dict, Optional, Union, Any
  6.  
  7. import structlog
  8.  
  9. import setup_logger
  10.  
  11. logger = structlog.get_logger(__name__)
  12.  
  13.  
  14. class MySignal:
  15.     def __init__(self):
  16.         self.handlers: List[Callable] = []
  17.  
  18.     def connect(self, handler: Callable):
  19.         self.handlers.append(handler)
  20.  
  21.     def emit(self, *args, **kwargs):
  22.         for handler in self.handlers:
  23.             handler(*args, **kwargs)
  24.  
  25.  
  26. class Handler:
  27.     def __init__(self):
  28.         self.signal = MySignal()
  29.  
  30.     def worker(self, name: str, timeout: int):
  31.         sleep(timeout)
  32.         self.signal.emit(name)
  33.  
  34.  
  35. async def async_generator():
  36.     # Имитация NATS Subscribe
  37.     for i, data in enumerate([
  38.         ("task1", 10), ("task2", 7), ("status", "task2"), ("status", "no_task"), ("task3", 6), ("status", "task3"),
  39.         ("task4", 3), ("status", "task2"), ("task5", 3),
  40.         ("exit", 1)], start=1
  41.     ):
  42.         if data[0] == "exit":
  43.             await asyncio.sleep(data[1])
  44.             continue
  45.  
  46.         await asyncio.sleep(1)
  47.         yield i, data
  48.  
  49.  
  50. class WorkerThread(threading.Thread):
  51.     def __init__(self, queue: 'Queue'):
  52.         threading.Thread.__init__(self)
  53.         self._exit = False
  54.         self.queue = queue
  55.         self.start()
  56.  
  57.     def run(self):
  58.         while not self._exit:
  59.             task: Optional[tuple[Any, ...]] = self.queue.get_task()
  60.             if task is not None:
  61.                 name, data = task
  62.                 self.queue.set_work(name, self)
  63.  
  64.                 handler, args, kwargs = data
  65.                 handler(*args, **kwargs)
  66.                 self.queue.done_work(name)
  67.  
  68.  
  69. class Status(StrEnum):
  70.     RUNNING = "running"
  71.     WAITING = "waiting"
  72.     DONE = "done"
  73.  
  74.  
  75. class Queue:
  76.     def __init__(self, max_thread: int = 0):
  77.         self._queue = {}
  78.         self.pool = {}
  79.         self._workers = []
  80.         self.init_threads(max_thread)
  81.  
  82.     def __enter__(self):
  83.         # Инициализация ресурсов или выполнение действий при входе в блок with
  84.         logger.debug("Run queue")
  85.         print()
  86.         return self
  87.  
  88.     def __exit__(self, exc_type, exc_value, traceback):
  89.         # Завершение работы, освобождение ресурсов или выполнение действий при выходе из блока with
  90.         logger.debug("Exit...")
  91.         # while not len(self._queue) == 0: ...
  92.         for worker in self._workers:
  93.             worker._exit = True
  94.             worker.join()
  95.         logger.debug("Exit queue")
  96.         # Здесь вы можете выполнить необходимые действия при завершении работы класса
  97.  
  98.     def init_threads(self, max_threads: int):
  99.         for _ in range(max_threads):
  100.             self._workers.append(WorkerThread(self))
  101.         logger.debug(f"Create {max_threads} threads")
  102.  
  103.     def get_task(self) -> Optional[tuple[Any, ...]]:
  104.         try:
  105.             return self._queue.popitem()
  106.         except KeyError:
  107.             return None
  108.  
  109.     def set_work(self, name: str, worker: 'WorkerThread'):
  110.         self.pool[name] = worker
  111.  
  112.     def done_work(self, name: str):
  113.         self.pool.pop(name, None)
  114.  
  115.     def get_status(self, name: str):
  116.         if self._queue.get(name, None) is not None:
  117.             return Status.WAITING
  118.         elif self.pool.get(name, None) is not None:
  119.             return Status.RUNNING
  120.         return Status.DONE
  121.  
  122.     def create_task(self, __handler: Callable, __name: str, *args, **kwargs):
  123.         self._queue[__name] = (__handler, args, kwargs)
  124.  
  125.  
  126. def output_handler(name: str):
  127.     logger.debug("<=== Work completed", name=name)
  128.  
  129.  
  130. async def main():
  131.     with Queue(2) as queue:
  132.         async for i, data in async_generator():
  133.             if data[0] == "status":
  134.                 name = data[1]
  135.                 status: Status = queue.get_status(name)
  136.                 logger.debug(f"Status {name}: {status}")
  137.             else:
  138.                 logger.debug("===> New task", name=data[0], timeout=data[1])
  139.  
  140.                 handler = Handler()
  141.                 handler.signal.connect(output_handler)
  142.                 queue.create_task(handler.worker, data[0], name=data[0], timeout=data[1])
  143.  
  144.  
  145. if __name__ == '__main__':
  146.     asyncio.run(main())
  147.  
  148.  
  149.  
  150. LOGS:
  151. 2024-01-24 14:28:54 [debug    ] Create 2 threads               [test_thread_3.py:init_threads:112]
  152. 2024-01-24 14:28:57 [debug    ] Run queue                      [test_thread_3.py:__enter__:95]
  153.  
  154. 2024-01-24 14:28:58 [debug    ] ===> New task                  [test_thread_3.py:main:149] name=task1 timeout=10
  155. 2024-01-24 14:29:00 [debug    ] ===> New task                  [test_thread_3.py:main:149] name=task2 timeout=7
  156. 2024-01-24 14:29:01 [debug    ] Status task2: running          [test_thread_3.py:main:147]
  157. 2024-01-24 14:29:02 [debug    ] Status no_task: done           [test_thread_3.py:main:147]
  158. 2024-01-24 14:29:03 [debug    ] ===> New task                  [test_thread_3.py:main:149] name=task3 timeout=6
  159. 2024-01-24 14:29:04 [debug    ] Status task3: waiting          [test_thread_3.py:main:147]
  160. 2024-01-24 14:29:05 [debug    ] ===> New task                  [test_thread_3.py:main:149] name=task4 timeout=3
  161. 2024-01-24 14:29:06 [debug    ] Status task2: running          [test_thread_3.py:main:147]
  162. 2024-01-24 14:29:07 [debug    ] <=== Work completed            [test_thread_3.py:output_handler:138] name=task2
  163. 2024-01-24 14:29:07 [debug    ] ===> New task                  [test_thread_3.py:main:149] name=task5 timeout=3
  164. 2024-01-24 14:29:08 [debug    ] Exit...                        [test_thread_3.py:__exit__:101]
  165. 2024-01-24 14:29:09 [debug    ] <=== Work completed            [test_thread_3.py:output_handler:138] name=task1
  166. 2024-01-24 14:29:10 [debug    ] <=== Work completed            [test_thread_3.py:output_handler:138] name=task4
  167. 2024-01-24 14:29:12 [debug    ] <=== Work completed            [test_thread_3.py:output_handler:138] name=task5
  168. 2024-01-24 14:29:12 [debug    ] Exit queue                     [test_thread_3.py:__exit__:106]
Advertisement
Add Comment
Please, Sign In to add comment