Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- class Thread(threading.Thread):
- def __init__(self, queue: 'Queue'):
- threading.Thread.__init__(self)
- self.queue = queue
- self.__exit = False
- self.start()
- def run(self) -> None:
- while not self.__exit:
- task: Optional[tuple[str, tuple]] = self.queue.get_task()
- if task is not None:
- name, data = task
- self.queue.took_task(name, self)
- handler, args, kwargs = data
- try:
- if handler.is_async:
- loop = asyncio.new_event_loop()
- asyncio.set_event_loop(loop)
- loop.run_until_complete(handler.worker(*args, **kwargs))
- loop.close()
- else:
- handler.worker(*args, **kwargs)
- except ErrorHandler as exc:
- ...
- finally:
- self.queue.done_task(name)
- class Queue:
- def __init__(self, max_thread: int, debug: bool = False) -> None:
- self.max_thread = max_thread
- self.__debug = debug
- self.__queue: Dict[str, tuple] = {}
- self.__pool: Dict[str, Thread] = {}
- self.__workers: List[Thread] = []
- self.workers_init()
- def workers_init(self):
- if not self.__debug:
- for _ in range(self.max_thread):
- self.__workers.append(Thread(queue=self))
- else:
- from tqdm import trange
- pbar = trange(self.max_thread)
- for numerate in pbar:
- pbar.set_description(f"Create thread #'{numerate}'")
- self.__workers.append(Thread(queue=self))
- def took_task(self, name: str, thread: 'Thread') -> None:
- self.__pool[name] = thread
- def done_task(self, name: str):
- self.__pool.pop(name, None)
- def add_task(self, __handler: Handler, __name: str, *args, **kwargs) -> None:
- self.__queue[__name] = (__handler, args, kwargs)
- def get_task(self) -> Optional[tuple[str, tuple]]:
- try:
- return self.__queue.popitem()
- except KeyError:
- return None
- def get_status(self, name: str) -> Optional[Status]:
- if name in self.__queue:
- return Status.WAITING
- elif name in self.__pool:
- return Status.RUNNING
- return Status.DONE
Advertisement
Add Comment
Please, Sign In to add comment