Not a member of Pastebin yet?
Sign Up,
it unlocks many cool features!
- import asyncio
- import structlog
- from time import time
- from dataclasses import dataclass
- from typing import Callable, Awaitable, Any
- logger = structlog.get_logger()
- @dataclass
- class TaskModel:
- name: str
- handler: Callable[..., Awaitable[None]]
- vars: dict[str, Any]
- complete_callback: Callable[..., Awaitable[None]] = None
- error_callback: Callable[..., Awaitable[None]] = None
- close_callback: Callable[..., Awaitable[None]] = None
- waiting_callback: Callable[..., Awaitable[None]] = None
- class Task:
- WAITING_PROCESS = 5
- PROCESS_TIMEOUT = 10
- is_ready = False
- def __init__(self, model: TaskModel) -> None:
- self.model = model
- async def wait_result(self):
- start_time = time()
- while not self.is_ready:
- await asyncio.sleep(0.2)
- if 0 < self.PROCESS_TIMEOUT < time() - start_time:
- raise TimeoutError
- async def run(self):
- try:
- # await self.wait_result()
- self.is_ready = await self.model.handler(name=self.model.name, **self.model.vars)
- if self.model.complete_callback:
- await self.model.complete_callback(name=self.model.name)
- except asyncio.CancelledError:
- if self.model.close_callback:
- await self.model.close_callback(name=self.model.name)
- except Exception as exc:
- if self.model.error_callback:
- await self.model.error_callback(name=self.model.name, exc=exc)
- class TaskManager:
- tasks = []
- async def create_task(
- self, name, test_timeout, handler,
- callback=None, error_callback=None, close_callback=None, waiting_callback=None
- ):
- new_task = Task(
- TaskModel(
- name=name,
- handler=handler,
- vars={"test_timeout": test_timeout},
- complete_callback=callback,
- error_callback=error_callback,
- close_callback=close_callback,
- waiting_callback=waiting_callback
- )
- )
- task = asyncio.create_task(new_task.run())
- logger.debug(f"Create {name}", test_timeout=test_timeout)
- self.tasks.append({
- "id": name,
- "task_object": new_task,
- "task": task
- })
- def close_task(self, name):
- for task in self.tasks:
- if task["id"] == name:
- task["task"].cancel()
- self.tasks.remove(task)
- logger.debug(f"Closed", name=name)
- break
- def get_status(self, name):
- for task in self.tasks:
- if task["id"] == name:
- return task["task_object"].is_ready
- raise ValueError
- async def callback_handler(name):
- logger.debug("Completed", name=name)
- async def close_handler(name):
- logger.warning(f"{name} Close")
- async def error_handler(name, exc):
- logger.warning(f"{name} Error", exc=exc)
- async def add_waiting(name):
- logger.info(f"{name} Waiting")
- async def test_handler(name, test_timeout):
- logger.debug(f"{name} Run task")
- await asyncio.sleep(test_timeout)
- if name == "Task2":
- raise ValueError("Test")
- return True
- async def async_generator():
- # Имитация NATS Subscribe
- for i, data in enumerate([
- ("task", 7), ("task", 2), ("task", 8), ("status", "Task3", 1), ("close", "Task3", 2),
- ("exit", 10)], start=1
- ):
- if data[0] == "exit":
- await asyncio.sleep(data[1])
- continue
- elif data[0] == "close":
- await asyncio.sleep(data[2])
- elif data[0] == "status":
- await asyncio.sleep(data[2])
- await asyncio.sleep(1)
- yield i, data
- async def main():
- task_manager = TaskManager()
- # while True:
- async for i, data in async_generator():
- _input = data[0]
- metadata = data[1]
- if _input == "task":
- await task_manager.create_task(
- name=f"Task{i}",
- handler=test_handler,
- test_timeout=metadata,
- callback=callback_handler,
- error_callback=error_handler,
- close_callback=close_handler,
- waiting_callback=add_waiting
- )
- elif _input == "close":
- task_manager.close_task(metadata)
- elif _input == "status":
- try:
- status = task_manager.get_status(metadata)
- logger.debug(f"Task {metadata}: {status}")
- except ValueError:
- logger.debug(f"Task {metadata} is not found")
- if __name__ == "__main__":
- asyncio.run(main())
Advertisement
Add Comment
Please, Sign In to add comment