morington

Untitled

Jan 26th, 2024
1,017
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
Python 2.42 KB | None | 0 0
  1. class Thread(threading.Thread):
  2.     def __init__(self, queue: 'Queue'):
  3.         threading.Thread.__init__(self)
  4.         self.queue = queue
  5.  
  6.         self.__exit = False
  7.         self.start()
  8.  
  9.     def run(self) -> None:
  10.         while not self.__exit:
  11.             task: Optional[tuple[str, tuple]] = self.queue.get_task()
  12.             if task is not None:
  13.                 name, data = task
  14.                 self.queue.took_task(name, self)
  15.  
  16.                 handler, args, kwargs = data
  17.                 try:
  18.                     if handler.is_async:
  19.                         loop = asyncio.new_event_loop()
  20.                         asyncio.set_event_loop(loop)
  21.                         loop.run_until_complete(handler.worker(*args, **kwargs))
  22.                         loop.close()
  23.                     else:
  24.                         handler.worker(*args, **kwargs)
  25.                 except ErrorHandler as exc:
  26.                     ...
  27.                 finally:
  28.                     self.queue.done_task(name)
  29.  
  30.  
  31. class Queue:
  32.     def __init__(self, max_thread: int, debug: bool = False) -> None:
  33.         self.max_thread = max_thread
  34.         self.__debug = debug
  35.  
  36.         self.__queue: Dict[str, tuple] = {}
  37.         self.__pool: Dict[str, Thread] = {}
  38.         self.__workers: List[Thread] = []
  39.         self.workers_init()
  40.  
  41.     def workers_init(self):
  42.         if not self.__debug:
  43.             for _ in range(self.max_thread):
  44.                 self.__workers.append(Thread(queue=self))
  45.         else:
  46.             from tqdm import trange
  47.  
  48.             pbar = trange(self.max_thread)
  49.             for numerate in pbar:
  50.                 pbar.set_description(f"Create thread #'{numerate}'")
  51.                 self.__workers.append(Thread(queue=self))
  52.  
  53.     def took_task(self, name: str, thread: 'Thread') -> None:
  54.         self.__pool[name] = thread
  55.  
  56.     def done_task(self, name: str):
  57.         self.__pool.pop(name, None)
  58.  
  59.     def add_task(self, __handler: Handler, __name: str, *args, **kwargs) -> None:
  60.         self.__queue[__name] = (__handler, args, kwargs)
  61.  
  62.     def get_task(self) -> Optional[tuple[str, tuple]]:
  63.         try:
  64.             return self.__queue.popitem()
  65.         except KeyError:
  66.             return None
  67.  
  68.     def get_status(self, name: str) -> Optional[Status]:
  69.         if name in self.__queue:
  70.             return Status.WAITING
  71.         elif name in self.__pool:
  72.             return Status.RUNNING
  73.         return Status.DONE
Advertisement
Add Comment
Please, Sign In to add comment