refactoring

This commit is contained in:
Maxim Devaev
2022-06-14 11:23:04 +03:00
parent 6caeb2ce82
commit e050bbd725
2 changed files with 50 additions and 34 deletions

View File

@@ -20,8 +20,6 @@
# ========================================================================== #
import os
import signal
import asyncio
import operator
import dataclasses
@@ -208,8 +206,6 @@ class KvmdServer(HttpServer): # pylint: disable=too-many-arguments,too-many-ins
self.__ws_clients: Set[_WsClient] = set()
self.__ws_clients_lock = asyncio.Lock()
self.__system_tasks: List[asyncio.Task] = []
self.__streamer_notifier = aiotools.AioNotifier()
self.__reset_streamer = False
self.__new_streamer_params: Dict = {}
@@ -297,13 +293,13 @@ class KvmdServer(HttpServer): # pylint: disable=too-many-arguments,too-many-ins
await check_request_auth(self.__auth_manager, exposed, request)
async def _init_app(self) -> None:
self.__run_system_task(self.__stream_controller)
aiotools.create_deadly_task("Stream controller", self.__stream_controller())
for comp in self.__components:
if comp.systask:
self.__run_system_task(comp.systask)
aiotools.create_deadly_task(comp.name, comp.systask())
if comp.poll_state:
self.__run_system_task(self.__poll_state, comp.event_type, comp.poll_state())
self.__run_system_task(self.__stream_snapshoter)
aiotools.create_deadly_task(f"{comp.name} [poller]", self.__poll_state(comp.event_type, comp.poll_state()))
aiotools.create_deadly_task("Stream snapshoter", self.__stream_snapshoter())
for api in self.__apis:
for http_exposed in get_exposed_http(api):
@@ -311,31 +307,14 @@ class KvmdServer(HttpServer): # pylint: disable=too-many-arguments,too-many-ins
for ws_exposed in get_exposed_ws(api):
self.__ws_handlers[ws_exposed.event_type] = ws_exposed.handler
def __run_system_task(self, method: Callable, *args: Any) -> None:
async def wrapper() -> None:
try:
await method(*args)
raise RuntimeError(f"Dead system task: {method}"
f"({', '.join(getattr(arg, '__name__', str(arg)) for arg in args)})")
except asyncio.CancelledError:
pass
except Exception:
get_logger().exception("Unhandled exception, killing myself ...")
os.kill(os.getpid(), signal.SIGTERM)
self.__system_tasks.append(asyncio.create_task(wrapper()))
async def _on_shutdown(self) -> None:
logger = get_logger(0)
logger.info("Waiting short tasks ...")
await asyncio.gather(*aiotools.get_short_tasks(), return_exceptions=True)
await aiotools.wait_all_short_tasks()
logger.info("Cancelling system tasks ...")
for task in self.__system_tasks:
task.cancel()
logger.info("Waiting system tasks ...")
await asyncio.gather(*self.__system_tasks, return_exceptions=True)
logger.info("Stopping system tasks ...")
await aiotools.stop_all_deadly_tasks()
logger.info("Disconnecting clients ...")
for client in list(self.__ws_clients):