14SIGNAL_ERRORS: dict[int, str] = {
16 'Service crashed with {signal_name} signal (segmentation fault)'
18 signal.SIGABRT:
'Service aborted by {signal_name} signal',
20DEFAULT_SIGNAL_ERROR =
'Service terminated by {signal_name} signal'
22_KNOWN_SIGNALS: dict[int, str] = {
23 signal.SIGABRT:
'SIGABRT',
24 signal.SIGBUS:
'SIGBUS',
25 signal.SIGFPE:
'SIGFPE',
26 signal.SIGHUP:
'SIGHUP',
27 signal.SIGINT:
'SIGINT',
28 signal.SIGKILL:
'SIGKILL',
29 signal.SIGPIPE:
'SIGPIPE',
30 signal.SIGSEGV:
'SIGSEGV',
31 signal.SIGTERM:
'SIGTERM',
35logger = logging.getLogger(__name__)
47 def __init__(self, message: str, exit_code: int) ->
None:
48 super().__init__(message)
49 self.exit_code = exit_code
54 def __init__(self, loop=None):
58 async def aclose(self):
60 await asyncio.wait(self.
_tasks)
62 async def add(self, pipe, handler):
63 if not pipe
or not handler:
65 is_coro = inspect.iscoroutinefunction(handler)
66 reader = await _create_pipe_reader(pipe, loop=self.
_loop)
68 async def data_handler():
69 async for line
in reader:
76 self.
_tasks.append(asyncio.create_task(coro))
79@contextlib.asynccontextmanager
83 shutdown_signal: int = signal.SIGINT,
84 shutdown_timeout: float = 120,
85 subprocess_spawner=
None,
89) -> AsyncGenerator[subprocess.Popen,
None]:
91 kwargs[
'stdout'] = subprocess.PIPE
93 kwargs[
'stderr'] = subprocess.PIPE
95 kwargs[
'preexec_fn'] = kwargs.get(
'preexec_fn', _setup_process)
97 logger.debug(
'Starting process with args %r', args)
98 if subprocess_spawner:
99 process = subprocess_spawner(args, **kwargs)
101 process = subprocess.Popen(args, **kwargs)
102 logging.debug(
'[%d] Process started', process.pid)
105 await readers.add(process.stdout, stdout_handler)
106 await readers.add(process.stderr, stderr_handler)
108 async with contextlib.aclosing(readers):
109 async with _shutdown_service(
111 shutdown_signal=shutdown_signal,
112 shutdown_timeout=shutdown_timeout,
117def exit_code_error(process: subprocess.Popen) -> ExitCodeError:
118 retcode = process.returncode
122def _exit_code_text(retcode: int):
124 return f
'Service exited with status code {retcode}'
125 signal_name = _pretty_signal(-retcode)
126 signal_error_fmt = SIGNAL_ERRORS.get(-retcode, DEFAULT_SIGNAL_ERROR)
127 return signal_error_fmt.format(signal_name=signal_name)
130@contextlib.asynccontextmanager
131async def _shutdown_service(*args, **kwargs):
135 await _do_service_shutdown(*args, **kwargs)
138async def _do_service_shutdown(process, *, shutdown_signal, shutdown_timeout):
139 allowed_exit_codes = (-shutdown_signal, 0)
141 retcode = process.poll()
142 if retcode
is not None:
144 '[%d] Process already finished with code %d', process.pid, retcode
146 if retcode
not in allowed_exit_codes:
147 raise exit_code_error(process)
151 process.send_signal(shutdown_signal)
156 '[%d] Trying to stop process with signal %s',
158 _pretty_signal(shutdown_signal),
160 poll_start = time.monotonic()
162 retcode = process.poll()
163 if retcode
is not None:
164 if retcode
not in allowed_exit_codes:
165 raise exit_code_error(process)
167 current_time = time.monotonic()
168 if current_time - poll_start > shutdown_timeout:
170 await asyncio.sleep(_POLL_TIMEOUT)
173 '[%d] Process did not finished within shutdown timeout %d seconds',
178 logger.warning(
'[%d] Now killing process with signal SIGKILL', process.pid)
180 retcode = process.poll()
181 if retcode
is not None:
182 raise exit_code_error(process)
184 process.send_signal(signal.SIGKILL)
187 await asyncio.sleep(_POLL_TIMEOUT)
190def _pretty_signal(signum: int) -> str:
191 if signum
in _KNOWN_SIGNALS:
192 return _KNOWN_SIGNALS[signum]
196async def _create_pipe_reader(pipe, loop=None):
198 loop = asyncio.get_running_loop()
199 reader = asyncio.StreamReader(loop=loop)
200 reader_protocol = asyncio.StreamReaderProtocol(reader)
201 await loop.connect_read_pipe(
lambda: reader_protocol, pipe)
207if sys.platform ==
'linux':
208 _LIBC = ctypes.CDLL(
'libc.so.6')
213def _setup_process() -> None:
214 if _LIBC
is not None:
215 _LIBC.prctl(_PR_SET_PDEATHSIG, signal.SIGKILL)
218__tracebackhide__ = traceback.hide(BaseError)