Исходный код 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()