Deep Engineering

MEASUREMENT

bench/shutdown/drain.py

The script that produced the numbers in the article, and the record of the run. The file is read from the repository at build time — this is the code that was run, not a copy of it.

Cited in
/en/interview/sre/graceful-shutdown
How to run it
python3 bench/shutdown/drain.py    > bench/shutdown/runs/drain.txt
python3 bench/shutdown/practice.py > bench/shutdown/runs/practice.txt

The run below is recorded in Russian. It is a lab record, kept in the language it was written in; the numbers, the tables and the code read the same either way.

Record of the run

Замеры для урока «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.

Script

192 lines
"""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()