morington

Untitled

Jan 23rd, 2024 (edited)
730
0
Never
Not a member of Pastebin yet? Sign Up, it unlocks many cool features!
Python 4.77 KB | None | 0 0
  1. import asyncio
  2. import structlog
  3. from time import time
  4. from dataclasses import dataclass
  5. from typing import Callable, Awaitable, Any
  6.  
  7. logger = structlog.get_logger()
  8.  
  9.  
  10. @dataclass
  11. class TaskModel:
  12.     name: str
  13.     handler: Callable[..., Awaitable[None]]
  14.     vars: dict[str, Any]
  15.     complete_callback: Callable[..., Awaitable[None]] = None
  16.     error_callback: Callable[..., Awaitable[None]] = None
  17.     close_callback: Callable[..., Awaitable[None]] = None
  18.     waiting_callback: Callable[..., Awaitable[None]] = None
  19.  
  20.  
  21. class Task:
  22.     WAITING_PROCESS = 5
  23.     PROCESS_TIMEOUT = 10
  24.     is_ready = False
  25.  
  26.     def __init__(self, model: TaskModel) -> None:
  27.         self.model = model
  28.  
  29.     async def wait_result(self):
  30.         start_time = time()
  31.         while not self.is_ready:
  32.             await asyncio.sleep(0.2)
  33.             if 0 < self.PROCESS_TIMEOUT < time() - start_time:
  34.                 raise TimeoutError
  35.  
  36.     async def run(self):
  37.         try:
  38.             # await self.wait_result()
  39.             self.is_ready = await self.model.handler(name=self.model.name, **self.model.vars)
  40.  
  41.             if self.model.complete_callback:
  42.                 await self.model.complete_callback(name=self.model.name)
  43.  
  44.         except asyncio.CancelledError:
  45.             if self.model.close_callback:
  46.                 await self.model.close_callback(name=self.model.name)
  47.  
  48.         except Exception as exc:
  49.             if self.model.error_callback:
  50.                 await self.model.error_callback(name=self.model.name, exc=exc)
  51.  
  52.  
  53. class TaskManager:
  54.     tasks = []
  55.  
  56.     async def create_task(
  57.             self, name, test_timeout, handler,
  58.             callback=None, error_callback=None, close_callback=None, waiting_callback=None
  59.     ):
  60.         new_task = Task(
  61.             TaskModel(
  62.                 name=name,
  63.                 handler=handler,
  64.                 vars={"test_timeout": test_timeout},
  65.                 complete_callback=callback,
  66.                 error_callback=error_callback,
  67.                 close_callback=close_callback,
  68.                 waiting_callback=waiting_callback
  69.             )
  70.         )
  71.  
  72.         task = asyncio.create_task(new_task.run())
  73.         logger.debug(f"Create {name}", test_timeout=test_timeout)
  74.         self.tasks.append({
  75.             "id": name,
  76.             "task_object": new_task,
  77.             "task": task
  78.         })
  79.  
  80.     def close_task(self, name):
  81.         for task in self.tasks:
  82.             if task["id"] == name:
  83.                 task["task"].cancel()
  84.                 self.tasks.remove(task)
  85.                 logger.debug(f"Closed", name=name)
  86.                 break
  87.  
  88.     def get_status(self, name):
  89.         for task in self.tasks:
  90.             if task["id"] == name:
  91.                 return task["task_object"].is_ready
  92.         raise ValueError
  93.  
  94.  
  95. async def callback_handler(name):
  96.     logger.debug("Completed", name=name)
  97.  
  98.  
  99. async def close_handler(name):
  100.     logger.warning(f"{name} Close")
  101.  
  102.  
  103. async def error_handler(name, exc):
  104.     logger.warning(f"{name} Error", exc=exc)
  105.  
  106.  
  107. async def add_waiting(name):
  108.     logger.info(f"{name} Waiting")
  109.  
  110.  
  111. async def test_handler(name, test_timeout):
  112.     logger.debug(f"{name} Run task")
  113.     await asyncio.sleep(test_timeout)
  114.     if name == "Task2":
  115.         raise ValueError("Test")
  116.  
  117.     return True
  118.  
  119.  
  120. async def async_generator():
  121.     # Имитация NATS Subscribe
  122.     for i, data in enumerate([
  123.         ("task", 7), ("task", 2), ("task", 8), ("status", "Task3", 1), ("close", "Task3", 2),
  124.         ("exit", 10)], start=1
  125.     ):
  126.         if data[0] == "exit":
  127.             await asyncio.sleep(data[1])
  128.             continue
  129.         elif data[0] == "close":
  130.             await asyncio.sleep(data[2])
  131.         elif data[0] == "status":
  132.             await asyncio.sleep(data[2])
  133.  
  134.         await asyncio.sleep(1)
  135.         yield i, data
  136.  
  137.  
  138. async def main():
  139.     task_manager = TaskManager()
  140.  
  141.     # while True:
  142.     async for i, data in async_generator():
  143.         _input = data[0]
  144.         metadata = data[1]
  145.  
  146.         if _input == "task":
  147.             await task_manager.create_task(
  148.                 name=f"Task{i}",
  149.                 handler=test_handler,
  150.                 test_timeout=metadata,
  151.                 callback=callback_handler,
  152.                 error_callback=error_handler,
  153.                 close_callback=close_handler,
  154.                 waiting_callback=add_waiting
  155.             )
  156.         elif _input == "close":
  157.             task_manager.close_task(metadata)
  158.         elif _input == "status":
  159.             try:
  160.                 status = task_manager.get_status(metadata)
  161.                 logger.debug(f"Task {metadata}: {status}")
  162.             except ValueError:
  163.                 logger.debug(f"Task {metadata} is not found")
  164.  
  165.  
  166. if __name__ == "__main__":
  167.     asyncio.run(main())
  168.  
Advertisement
Add Comment
Please, Sign In to add comment