Исходный код soniks_client.subprocess_log

"""Запуск подпроцесса с чтением его вывода в лог отдельным потоком.

Общее для потокового графа и скриптов вокруг прохода: одни и те же параметры
``Popen`` (расхождение в ``bufsize``, ``encoding`` или ``errors`` молча ломает
чтение вывода) и один и тот же слив ``stdout``.

Читать вывод именно потоком обязательно: ``for line in process.stdout``
в основном блоке висит до EOF, из-за чего ``wait(timeout)`` становится
недостижим и зависший процесс держит проход навсегда.
"""

import subprocess
from collections.abc import Callable
from threading import Thread

from core.configs import logger

# Сколько ждать поток чтения вывода после того, как процесс уже завершился.
JOIN_TIMEOUT_IN_SECONDS = 5


[документация] def start_logged_process( argv: list[str], log: Callable[[str], None], env: dict[str, str] | None = None, ) -> tuple[subprocess.Popen, Thread]: """Запустить процесс и поток, сливающий его вывод в лог. Args: argv: Команда целиком, первым элементом — исполняемый файл. env: Окружение процесса целиком; ``None`` — унаследовать своё. log: Куда писать очередную строку вывода. Уровень и префикс — забота вызывающего, здесь про них ничего не известно. Raises: OSError: Процесс не удалось запустить. Гасить решает вызывающий: у скрипта прохода и у потокового графа политика разная. """ process = subprocess.Popen( argv, stdout=subprocess.PIPE, stderr=subprocess.STDOUT, text=True, bufsize=1, encoding="utf-8", errors="replace", env=env, ) log_thread = Thread(target=_drain, args=(process, log), daemon=True) log_thread.start() return process, log_thread
[документация] def stop_logging(log_thread: Thread | None, name: str) -> None: """Дождаться, пока поток чтения допишет остаток вывода. ``stdout`` закрывает сам поток чтения, а не вызывающий: ``close()`` ждёт лок буфера, который читатель держит всё время чтения. Скрипт, оставивший после себя фоновый процесс (``bandscan.sh start`` из ``satnogs-post``), отдаёт ему свой ``stdout`` по наследству — и закрытие отсюда блокировало вызывающего до смерти этого процесса, то есть до следующего прохода. Таймаут на процесс от этого не спасал: сам скрипт к тому моменту давно завершился. """ if log_thread is None: return log_thread.join(timeout=JOIN_TIMEOUT_IN_SECONDS) if log_thread.is_alive(): logger.warning( "Вывод %s ещё держат фоновые процессы, чтение продолжится в фоне", name, )
def _drain(process: subprocess.Popen, log: Callable[[str], None]) -> None: if process.stdout is None: return try: for line in process.stdout: if line.strip(): log(line.strip()) except ValueError: # Поток закрыт после принудительного завершения процесса. return finally: process.stdout.close()