From: Aleš Mrázek Date: Mon, 20 Jul 2026 13:46:58 +0000 (+0200) Subject: controller: supervisord controller X-Git-Url: http://git.ipfire.org/?a=commitdiff_plain;h=refs%2Fheads%2Fpython-refactoring-controller;p=thirdparty%2Fknot-resolver.git controller: supervisord controller --- diff --git a/python/knot_resolver/config/templates/supervisord.toml.j2 b/python/knot_resolver/config/templates/supervisord.toml.j2 index e69de29bb..055a71acd 100644 --- a/python/knot_resolver/config/templates/supervisord.toml.j2 +++ b/python/knot_resolver/config/templates/supervisord.toml.j2 @@ -0,0 +1,79 @@ +[supervisord] +directory = {{ supervisord.workdir }} +pidfile = {{ supervisord.pidfile }} +logfile = {{ supervisord.logfile }} +logfile_maxbytes = 0 +loglevel = {{ supervisord.loglevel }} +nodaemon = true +silent = true + +[supervisorctl] +serverurl = unix://{{ supervisord.unix_http_server }} + +[unix_http_server] +file = {{ supervisord.unix_http_server }} + +[rpcinterface:logger] +supervisor.rpcinterface_factory = knot_resolver.controller.plugins.logger_patch:inject +logtarget = {{ supervisord.logtarget }} + +[rpcinterface:notify] +supervisor.rpcinterface_factory = knot_resolver.controller.notify.notify_dispatcher:inject + +[program:manager] +command = {{ manager.command }} +directory = {{ manager.workdir }} +environment = {{ manager.environment }} +startsecs= {{ manager.startsecs }} +redirect_stderr = false +autostart = true +autorestart = true +killasgroup = true +stopsignal = SIGTERM +stopwaitsecs = 30 +stdout_logfile = NONE +stderr_logfile = NONE + +[program:worker] +process_name = %(program_name)s%(process_num)d +numprocs = {{ worker.max_procs }} +command = {{ worker.command }} +directory = {{ worker.workdir }} +environment = {{ worker.environment }} +startsecs = {{ worker.startsecs }} +redirect_stderr = false +autostart = false +autorestart = true +killasgroup = true +stopsignal = SIGTERM +stopwaitsecs = 30 +stdout_logfile = NONE +stderr_logfile = NONE + +[program:loader] +command = {{ loader.command }} +directory = {{ loader.workdir }} +environment = {{ loader.environment }} +startsecs = {{ loader.startsecs }} +redirect_stderr = false +autostart = false +exitcodes = 0 +killasgroup = true +stopsignal = SIGTERM +stopwaitsecs = 30 +stdout_logfile = NONE +stderr_logfile = NONE + +[program:cache-gc] +command = {{ cache_gc.command }} +directory = {{ cache_gc.workdir }} +environment = {{ cache_gc.environment }} +startsecs={{ cache_gc.startsecs }} +redirect_stderr=false +autostart=false +autorestart=true +killasgroup=true +stopsignal = SIGTERM +stopwaitsecs = 30 +stdout_logfile=NONE +stderr_logfile=NONE diff --git a/python/knot_resolver/controller/__init__.py b/python/knot_resolver/controller/__init__.py index e69de29bb..1343211e8 100644 --- a/python/knot_resolver/controller/__init__.py +++ b/python/knot_resolver/controller/__init__.py @@ -0,0 +1,5 @@ +from .supervisord import SupervisordController + +__all__ = [ + "SupervisordController", +] diff --git a/python/knot_resolver/controller/config.py b/python/knot_resolver/controller/config.py new file mode 100644 index 000000000..b882ab97a --- /dev/null +++ b/python/knot_resolver/controller/config.py @@ -0,0 +1,181 @@ +from __future__ import annotations + +import shutil +from dataclasses import dataclass +from pathlib import Path +from typing import TYPE_CHECKING, Literal + +from knot_resolver.config import KresConfig +from knot_resolver.constants import ( + CACHE_DIR, + KRES_CACHE_GC_EXECUTABLE, + KRES_MANAGER_EXECUTABLE, + KRESD_EXECUTABLE, + NOTIFY_SUPPORT, +) +from knot_resolver.logging import NO_PREFIX_FORMAT_ENV_VAR, get_logger + +if TYPE_CHECKING: + from knot_resolver.args import KresArgs + + SupervisordLogLevel = Literal["critical", "error", "warn", "info", "debug", "trace", "blather"] + SupervisordLogTarget = Literal["stdout", "stderr", "syslog"] + +# files names +SUPERVISORD_SOCK_NAME = "supervisord.sock" +SUPERVISORD_PIDFILE_NAME = "supervisord.pid" +SUPERVISORD_CONFIGFILE_NAME = "supervisord.conf" +SUPERVISORD_CONFIGFILE_NAME_TMP = f"{SUPERVISORD_CONFIGFILE_NAME}.tmp" +WORKER_CONFIGFILE_NAME = "worker%(process_num)d.conf" +LOADER_CONFIGFILE_NAME = "loader.conf" + +# logfiles/logdirs names +SUPERVISORD_LOGDIR_NAME = "logs" +SUPERVISORD_LOGFILE_NAME = "supervisord.log" +MANAGER_LOGFILE_NAME = "manager.log" +WORKER_LOGFILE_NAME = "worker%(process_num)d.log" +LOADER_LOGFILE_NAME = "loader.log" +CACHE_GC_LOGFILE_NAME = "cache_gc.log" + +X_TYPE_VAR_NAME = "X-SUPERVISORD-TYPE" +X_TYPE_NOTIFY = "notify" +INSTANCE_VAR_NAME = "SYSTEMD_INSTANCE" +ENVIRONMENT_TYPE_NOTIFY = f"{X_TYPE_VAR_NAME}={X_TYPE_NOTIFY}" +ENVIRONMENT_INSTANCE_NUM = f'{INSTANCE_VAR_NAME}="%(process_num)d"' + +logger = get_logger(__name__) + + +def get_absolute_path(path: str | Path | None = None) -> Path: + return Path(path).absolute() if path else Path().absolute() + + +def get_workdir_path() -> Path: + return get_absolute_path() + + +def get_logfile_path(logfile: str) -> Path: + return get_absolute_path(SUPERVISORD_LOGDIR_NAME) / logfile + + +@dataclass +class SupervisordConfig: + workdir: Path + pidfile: Path + logfile: Path + loglevel: SupervisordLogLevel + logtarget: SupervisordLogTarget + unix_http_server: Path + notify_support: bool = NOTIFY_SUPPORT + + @staticmethod + def create(args: KresArgs, config: KresConfig) -> SupervisordConfig: + return SupervisordConfig( + workdir=get_workdir_path(), + pidfile=get_absolute_path(SUPERVISORD_PIDFILE_NAME), + logfile=get_logfile_path(SUPERVISORD_LOGFILE_NAME), + loglevel="debug", # TODO(amrazek): use logging level from KresConfig + logtarget="stdout", # TODO(amrazek): use logging target from KresConfig + unix_http_server=get_absolute_path( + SUPERVISORD_SOCK_NAME # TODO(amrazek): use management API from KresConfig + ), + ) + + +@dataclass +class SubprocessConfig: + command: str + workdir: Path + logfile: Path + startsecs: int = 0 + max_procs: int = 1 + environment: str = "" + + @staticmethod + def create_manager(args: KresArgs, config: KresConfig) -> SubprocessConfig: + startsecs = 0 + environment = f"{NO_PREFIX_FORMAT_ENV_VAR}=true" + + if NOTIFY_SUPPORT: + startsecs = 600 + environment += f",{ENVIRONMENT_TYPE_NOTIFY}" + + if KRES_CACHE_GC_EXECUTABLE.exists(): + command_args = [str(KRES_MANAGER_EXECUTABLE)] + if not KRES_MANAGER_EXECUTABLE.exists(): + command_args = [ + str(shutil.which("python3")), + "-m", + "knot_resolver.manager", + ] + + command_args += [ + "--logtarget", + args.logtarget, + "--loglevel", + args.loglevel, + "--config", + *args.config, + ] + + return SubprocessConfig( + command=" ".join(command_args), + workdir=get_workdir_path(), + environment=environment, + startsecs=startsecs, + logfile=get_logfile_path(MANAGER_LOGFILE_NAME), + ) + + @staticmethod + def create_worker(args: KresArgs, config: KresConfig) -> SubprocessConfig: + max_procs = 1 + + # Default for non-Linux systems without support for systemd NOTIFY message. + # Therefore, we need to give the kresd workers a few seconds to start properly. + environment = ENVIRONMENT_INSTANCE_NUM + startsecs = 3 + + if NOTIFY_SUPPORT: + # There is support for systemd NOTIFY message. + # Here, 'startsecs' serves as a timeout for waiting for NOTIFY message. + environment += f",{ENVIRONMENT_TYPE_NOTIFY}" + startsecs = 60 + + config_path = get_absolute_path(WORKER_CONFIGFILE_NAME) + command_args: list[str] = [str(KRESD_EXECUTABLE), "--config", str(config_path), "-n"] + + return SubprocessConfig( + command=" ".join(command_args), + workdir=get_workdir_path(), + environment=environment, + startsecs=startsecs, + logfile=get_logfile_path(WORKER_LOGFILE_NAME), + max_procs=max_procs, + ) + + @staticmethod + def create_loader(args: KresArgs, config: KresConfig) -> SubprocessConfig: + config_path = get_absolute_path(LOADER_CONFIGFILE_NAME) + + command_args: list[str] = [str(KRESD_EXECUTABLE), "--config", str(config_path), "-c", "-", "-n"] + + return SubprocessConfig( + command=" ".join(command_args), + workdir=get_workdir_path(), + logfile=get_logfile_path(LOADER_LOGFILE_NAME), + ) + + @staticmethod + def create_cache_gc(args: KresArgs, config: KresConfig) -> SubprocessConfig: + # TODO(amrazek): use cache_dir from KresConfig + cache_dir = CACHE_DIR + + command_args: list[str] = [str(KRES_CACHE_GC_EXECUTABLE), "-c", str(cache_dir)] + + # TODO(amrazek): convert KresConfig to cache gc flags + + return SubprocessConfig( + command=" ".join(command_args), + workdir=get_workdir_path(), + logfile=get_logfile_path(CACHE_GC_LOGFILE_NAME), + ) diff --git a/python/knot_resolver/controller/supervisord.py b/python/knot_resolver/controller/supervisord.py new file mode 100644 index 000000000..52af460cf --- /dev/null +++ b/python/knot_resolver/controller/supervisord.py @@ -0,0 +1,104 @@ +# mypy: disable-error-code=import-untyped + +from __future__ import annotations + +import logging +import os +import shutil +from pathlib import Path +from typing import TYPE_CHECKING, cast +from xmlrpc.client import ServerProxy + +from supervisor.xmlrpc import SupervisorTransport + +from knot_resolver.config import KresConfig +from knot_resolver.config.templates import SUPERVISORD_TEMPLATE +from knot_resolver.logging import get_logger + +from .config import ( + SUPERVISORD_CONFIGFILE_NAME, + SUPERVISORD_CONFIGFILE_NAME_TMP, + SUPERVISORD_SOCK_NAME, + SubprocessConfig, + SupervisordConfig, + get_absolute_path, +) +from .errors import ControllerError + +if TYPE_CHECKING: + from typing import Any, Protocol + + from knot_resolver.args import KresArgs + + class SupervisorRPC(Protocol): + def getState(self) -> dict[str, Any]: ... + def getAllProcessInfo(self) -> list[dict[str, Any]]: ... + def startProcess(self, name: str, wait: bool = True) -> bool: ... + def stopProcess(self, name: str, wait: bool = True) -> bool: ... + + +logger = get_logger(__name__) + + +def _create_server_proxy(config: KresConfig) -> ServerProxy: + serverurl = "unix://" + str(get_absolute_path(SUPERVISORD_SOCK_NAME)) + transport = SupervisorTransport( + username=None, + password=None, + serverurl=serverurl, + ) + return ServerProxy("http://127.0.0.1", transport=transport) + + +def _create_supervisord_proxy(config: KresConfig) -> SupervisorRPC: + server_proxy = _create_server_proxy(config) + return cast(SupervisorRPC, server_proxy.supervisor) + + +class SupervisordController: + def __init__(self, args: KresArgs) -> None: + # TODO(amrazek): add declarative configuration + # self._config = config + self._args = args + + def write_config(self) -> None: + logger.notice("Creating supervisord controller configuration...") + + config: str = SUPERVISORD_TEMPLATE.render( + supervisord=SupervisordConfig.create(self._args), + manager=SubprocessConfig.create_manager(self._args), + worker=SubprocessConfig.create_worker(self._args), + loader=SubprocessConfig.create_loader(self._args), + cache_gc=SubprocessConfig.create_cache_gc(self._args), + ) + + config_path_tmp = Path(SUPERVISORD_CONFIGFILE_NAME_TMP) + with config_path_tmp.open("w") as file: + file.write(config) + config_path_tmp.rename(SUPERVISORD_CONFIGFILE_NAME) + + def exec(self) -> None: + supervisord = shutil.which("supervisord") + if not supervisord: + msg = "failed to find 'supervisord' executable" + raise ControllerError(msg) + + config_path = Path(SUPERVISORD_CONFIGFILE_NAME) + if not config_path.exists(): + msg = f"failed to find supervisord configuration file '{config_path}'" + raise ControllerError(msg) + + args = [ + str(supervisord), + "--configuration", + str(config_path), + ] + + logger.notice("Execing supervisord...") + logging.shutdown() + + try: + os.execv(supervisord, args) + except OSError as e: + msg = f"supervisord exec failed: {e}" + raise ControllerError(msg) from e