Deep Engineering

ЗАМЕР

bench/shutdown/drain.py

Скрипт, которым получены числа в статье, и запись прогона. Файл читается на сборке из репозитория — это тот самый код, который запускали, а не его копия.

Цитируется в статье
/ru/interview/sre/graceful-shutdown
Как запустить
python3 bench/shutdown/drain.py    > bench/shutdown/runs/drain.txt
python3 bench/shutdown/practice.py > bench/shutdown/runs/practice.txt

Запись прогона

Замеры для урока «Graceful shutdown»

Файл Что делает
drain.py настоящий сервер на loopback под нагрузкой из двадцати клиентов: сколько запросов обрывается при немедленном выходе по SIGTERM и при дренаже, и сколько занимает сама остановка
practice.py ответы к задачам урока: число оборванных запросов в обоих режимах, код выхода и время остановки

Запуск из корня репозитория:

python3 bench/shutdown/drain.py    > bench/shutdown/runs/drain.txt
python3 bench/shutdown/practice.py > bench/shutdown/runs/practice.txt

Что воспроизводимо

Соотношение: без дренажа обрываются все запросы, что были в работе; с дренажем — ни одного. И то, что цена дренажа равна остатку уже начатой работы, а не льготному сроку.

Не воспроизводимо: миллисекунды остановки и выдержка обработки — её задаёт сам скрипт (300 мс).

Требования к среде

Только loopback. Сервер живёт в отдельном процессе, чтобы ему можно было послать сигнал так, как это делает оркестратор.

Числа сняты на CPython 3.11.15, Linux 6.18.44.

Скрипт

192 строк
"""Graceful shutdown: что теряется в момент остановки и сколько это стоит.

ЗАЧЕМ ЭТОТ ФАЙЛ. «При каждой выкатке немного пятисоток» — фраза, которую
произносят как описание погоды. Здесь показано, что это не погода, а
следствие ровно одного решения: что делает процесс, получив `SIGTERM`, — выходит
сразу или доводит начатое.

ЧТО ЗДЕСЬ ИЗМЕРЯЕТСЯ. Число оборванных запросов при остановке под нагрузкой в
двух режимах и время самой остановки. Сервер — настоящий, на loopback, клиенты
— настоящие потоки; выдержка обработки задаётся скриптом.

ПОЧЕМУ ПОДПИСИ ПО-АНГЛИЙСКИ. Урок существует в двух языках и цитирует запись
прогона дословно обеими версиями.

ЗАПУСК: python3 bench/shutdown/drain.py
Вывод: runs/drain.txt
"""

import os
import signal
import socket
import subprocess
import sys
import threading
import time

HOST = "127.0.0.1"
WORK_MS = 300
CLIENTS = 20

# Сервер живёт в отдельном процессе — иначе ему нельзя послать сигнал так, как
# это делает оркестратор. Режим задаётся аргументом: `abrupt` выходит по
# SIGTERM немедленно, `drain` перестаёт принимать новые соединения и доводит
# начатые до конца.
SERVER = r"""
import os, signal, socket, sys, threading, time

mode, work_ms = sys.argv[1], int(sys.argv[2])
srv = socket.socket()
srv.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
srv.bind(("127.0.0.1", 0))
srv.listen(128)
print(srv.getsockname()[1], flush=True)

stopping = threading.Event()
inflight = threading.Semaphore(0)
active = 0
lock = threading.Lock()


def handle(conn):
    global active
    with lock:
        active += 1
    try:
        conn.recv(64)
        time.sleep(work_ms / 1000)
        conn.sendall(b"ok\n")
    except OSError:
        pass
    finally:
        conn.close()
        with lock:
            active -= 1


def on_term(signum, frame):
    if mode == "abrupt":
        # ПОЧЕМУ НЕ os._exit(143). Число 143 — соглашение оболочки (128 + номер
        # сигнала), а не то, что делает ядро. Назначить его самим значило бы
        # напечатать собственную константу под видом наблюдения. Поэтому здесь
        # обработчик снимается и сигнал посылается заново: процесс умирает от
        # SIGTERM по-настоящему, а родитель видит это в статусе ожидания.
        signal.signal(signal.SIGTERM, signal.SIG_DFL)
        os.kill(os.getpid(), signal.SIGTERM)
    stopping.set()


signal.signal(signal.SIGTERM, on_term)

def serve():
    while not stopping.is_set():
        try:
            srv.settimeout(0.05)
            conn, _ = srv.accept()
        except socket.timeout:
            continue
        except OSError:
            return
        threading.Thread(target=handle, args=(conn,), daemon=True).start()

worker = threading.Thread(target=serve, daemon=True)
worker.start()
stopping.wait()
srv.close()
# Дренаж: ждём, пока начатые запросы закончатся.
while True:
    with lock:
        if active == 0:
            break
    time.sleep(0.01)
os._exit(0)
"""


def show(title: str) -> None:
    print()
    print(title)
    print("-" * len(title))


def row(label: str, value: object) -> None:
    print(f"  {label:<44} {value}")


def start_server(mode: str) -> tuple[subprocess.Popen[str], int]:
    proc = subprocess.Popen(
        [sys.executable, "-c", SERVER, mode, str(WORK_MS)],
        stdout=subprocess.PIPE,
        text=True,
    )
    port = int(proc.stdout.readline().strip())
    return proc, port


def hammer(port: int, results: list[str], index: int) -> None:
    client = socket.socket()
    client.settimeout(5.0)
    try:
        client.connect((HOST, port))
        client.sendall(b"go\n")
        answer = client.recv(64)
        results[index] = "ok" if answer.strip() == b"ok" else "empty answer"
    except OSError as exc:
        results[index] = type(exc).__name__
    finally:
        client.close()


def run(mode: str) -> dict[str, object]:
    proc, port = start_server(mode)
    results = ["" for _ in range(CLIENTS)]
    threads = [
        threading.Thread(target=hammer, args=(port, results, i)) for i in range(CLIENTS)
    ]
    for thread in threads:
        thread.start()

    # Сигнал приходит В СЕРЕДИНЕ обработки: половина времени работы прошла.
    time.sleep(WORK_MS / 2000)
    started = time.perf_counter()
    proc.send_signal(signal.SIGTERM)
    exit_code = proc.wait(timeout=10)
    stop_ms = (time.perf_counter() - started) * 1000

    for thread in threads:
        thread.join()
    broken = sum(1 for r in results if r != "ok")
    return {"broken": broken, "stop_ms": stop_ms, "exit": exit_code}


def main() -> None:
    print(f"Python {sys.version.split()[0]} · Linux {os.uname().release} · loopback")
    print(f"{CLIENTS} concurrent requests, {WORK_MS} ms of work each")
    print("SIGTERM arrives halfway through the work")

    show("1. EXIT ON SIGTERM, THE WAY A PROCESS WITHOUT A HANDLER DOES")
    abrupt = run("abrupt")
    row("requests broken", f"{abrupt['broken']} of {CLIENTS}")
    row("time from SIGTERM to exit, ms", f"{abrupt['stop_ms']:.0f}")
    row("wait status: killed by signal", -int(abrupt["exit"]))
    row("the same as a shell reports it", 128 - int(abrupt["exit"]))

    show("2. STOP ACCEPTING, FINISH WHAT IS STARTED, THEN EXIT")
    drained = run("drain")
    row("requests broken", f"{drained['broken']} of {CLIENTS}")
    row("time from SIGTERM to exit, ms", f"{drained['stop_ms']:.0f}")
    row("exit code", drained["exit"])

    show("SIDE BY SIDE")
    row("broken requests, abrupt against drain", f"{abrupt['broken']} -> {drained['broken']}")
    row("stop time, abrupt against drain, ms",
        f"{abrupt['stop_ms']:.0f} -> {drained['stop_ms']:.0f}")
    print()
    print("  The whole difference is one decision: does the process leave at")
    print("  once or finish what it started. Draining costs the remainder of")
    print("  the work in flight - here about half of one request's time.")


if __name__ == "__main__":
    main()