Главная/Проекты/Решение
Эталонное решениеСначала выполните проект самостоятельно и используйте этот материал для сверки решений.

Проект 3. Эталонное решение

Открывать после собственной реализации и прохождения CHECKLIST.md.

Решение — не единственно верное. Места, где обоснован другой выбор, отмечены явно.

Что проверено на самом деле

ЧтоКакРезультат
Миграцииприменение к настоящему PostgreSQL 163 применены, повтор ничего не делает
Отказ при правке применённой миграцииизменение файла и повторный запускотказ с указанием обеих контрольных сумм
Хранилище13 тестов против настоящей базывсе проходят
Белый список полейпопытки подстановки в GROUP BYотклонены
Обработчик10 тестов, включая две копии сразуSKIP LOCKED работает
Три состояния тестовс базой, без базы, без базы с требованиемразличимы
Весь наборpytest37 пройдено, покрытие 80 %
Redis6 тестов против настоящего Redis 8все проходят
Стек целикомdocker compose up -dпять сервисов, четыре healthy, migrateExited (0)
Воспроизводимость стартадесять запусков подряддесять из десяти healthy
internal: trueпопытка выхода наружу из worker и из apiу worker выхода нет, у api есть — он ещё и во frontend
Состояниеdown / up и down -v / upпереживает down, исчезает при down -v, схема создаётся заново
Секретпоиск значения в конфигурации и в аргументах PID 1не найдено ни там, ни там

Прогон целиком — verify/62-project3.sh, 35 подтверждений из 35 (Docker Engine 29.7.1, 2026-08-05). Скрипт собирает стек из блоков этого файла и SOLUTION.md проекта 2: проверяется напечатанное, а не отдельная копия.

Первый запуск нашёл шесть дефектов, и первые три не давали стеку подняться вовсе.

ЧтоКак проявилось
PGPASSWORD_FILE не читал никтоfe_sendauth: no password supplied, migrate выходит с кодом 1
В Settings не было database_url и redis_urlUnknownSettingError, api в перезапусках
create_app не читал secretPoolTimeout: pool initialization incomplete after 30 sec
Стадия test прогоняла тесты, которым нужна базасборка падала: покрытие 20 % из одних пропусков
--cov=app вместо нового кода40 % и провал порога: файлы проекта 2 считались без своих тестов
У сервиса tests не было EVENTAPI_REQUIRE_DB=1недоступная база дала бы «33 skipped» и зелёный прогон

Первые три — одного происхождения: соглашение <ИМЯ>_FILE было принято за общий механизм Docker. Оно работает у образа postgres, потому что его entrypoint это реализует; для своего кода реализовать нужно самому (шаг 1a).


Шаг 1a. Пароль из secret

app/secrets.py:

python
"""Чтение секретов из файлов.

PGPASSWORD_FILE — не переменная libpq. libpq знает PGPASSWORD (значение)
и PGPASSFILE (файл формата .pgpass, с полями host:port:db:user:password).
Соглашение «<ИМЯ>_FILE» поддерживает entrypoint образа postgres — но
только для себя. Клиенту приходится читать файл самому.
"""
from __future__ import annotations

import os
from pathlib import Path


def read_secret(name: str, default: str | None = None) -> str | None:
    """Значение из <name>_FILE, затем из <name>, затем default.

    Приоритет у файла: secret надёжнее переменной окружения, которая
    видна в docker inspect и наследуется дочерними процессами.
    """
    path = os.environ.get(f"{name}_FILE")
    if path:
        try:
            return Path(path).read_text(encoding="utf-8").strip()
        except OSError as exc:
            raise RuntimeError(f"{name}_FILE={path}: {exc}") from exc
    return os.environ.get(name, default)


def export_pgpassword() -> None:
    """Положить пароль в PGPASSWORD, если он задан файлом.

    Вызывается в начале каждой точки входа. Пароль не подставляется
    в строку подключения намеренно: DSN попадает в аргументы процесса
    и виден в `ps` и в `docker inspect`, а переменная окружения — нет.
    """
    if "PGPASSWORD" in os.environ:
        return
    password = read_secret("PGPASSWORD")
    if password:
        os.environ["PGPASSWORD"] = password

Находка, ради которой этот шаг появился. Первая редакция передавала PGPASSWORD_FILE: /run/secrets/db_password в окружение всех трёх ролей и на этом останавливалась. Стек не поднимался вовсе:

console
$ docker compose up -d
 ✔ Container eventstack-db-1        Healthy
 ✘ Container eventstack-migrate-1   Error
service "migrate" didn't complete successfully: exit 1

$ docker compose logs migrate | tail -1
миграции: база недоступна за 30.0 с: connection failed:
connection to server at "172.21.0.2", port 5432 failed: fe_sendauth: no password supplied

Переменную PGPASSWORD_FILE не читал никто: ни libpq, ни код. Она была скопирована из настройки самого PostgreSQL (POSTGRES_PASSWORD_FILE), где такое соглашение действительно работает — но там его реализует entrypoint образа.

Отсюда правило: соглашение <ИМЯ>_FILE — не общий механизм Docker, а договорённость конкретного образа. Для своего кода его нужно реализовать самому (урок 11.5).


Шаг 1. Миграции

Три файла схемы:

migrations/001_events.sql:

sql
-- Схема событий.
CREATE TABLE events (
    id          BIGSERIAL PRIMARY KEY,
    path        TEXT             NOT NULL CHECK (path LIKE '/%'),
    method      TEXT             NOT NULL,
    status      INTEGER          NOT NULL CHECK (status BETWEEN 100 AND 599),
    duration_ms DOUBLE PRECISION NOT NULL CHECK (duration_ms >= 0),
    ts          TIMESTAMPTZ      NOT NULL DEFAULT now()
);

-- Индекс под самый частый запрос: выборка по времени в обратном порядке.
CREATE INDEX events_ts_desc ON events (ts DESC);

migrations/002_stats_indexes.sql:

sql
-- Индексы под группировку. Отдельной миграцией, потому что первая
-- уже применена в эксплуатации и переписывать её нельзя.
CREATE INDEX events_status ON events (status);
CREATE INDEX events_method ON events (method);
CREATE INDEX events_path ON events (path);

migrations/003_processed_flag.sql:

sql
-- Признак обработки события фоновым обработчиком.
ALTER TABLE events ADD COLUMN processed_at TIMESTAMPTZ;

-- Частичный индекс: строк с NULL обычно мало, и запрос обработчика
-- «дай необработанные» должен читать только их.
CREATE INDEX events_unprocessed ON events (id) WHERE processed_at IS NULL;

Почему индексы отдельной миграцией. Первая уже применена в эксплуатации, и переписывать её нельзя: у тех, кто применил, изменения не появятся. Это не формальность — именно поэтому раннер отвергает изменённый файл.

migrations/runner.py:

python
"""Применение миграций.

Три свойства, каждое из которых проверено тестом:

1. Порядок задаётся именем файла и не зависит от порядка обхода каталога.
2. Повторный запуск ничего не делает: применённое записано в таблицу.
3. Два экземпляра, запущенных одновременно, не мешают друг другу —
   консультативная блокировка PostgreSQL пропускает одного.

Третье свойство существеннее прочих: в Compose и в оркестраторе
сервис миграций может оказаться запущенным в нескольких копиях,
и без блокировки они выполнят одну и ту же миграцию дважды.
"""
from __future__ import annotations

import hashlib
import logging
import re
import sys
import time
from pathlib import Path

import psycopg

from app.secrets import export_pgpassword

log = logging.getLogger("migrations")

# Произвольное, но постоянное число: две копии сервиса должны выбрать
# один и тот же ключ, иначе блокировка ничего не даёт.
ADVISORY_LOCK_KEY = 0x6D6967_31

NAME_RE = re.compile(r"^(\d{3})_[a-z0-9_]+\.sql$")

BOOTSTRAP = """
CREATE TABLE IF NOT EXISTS schema_migrations (
    version     TEXT        PRIMARY KEY,
    checksum    TEXT        NOT NULL,
    applied_at  TIMESTAMPTZ NOT NULL DEFAULT now()
)
"""


class MigrationError(RuntimeError):
    pass


def discover(directory: Path) -> list[Path]:
    """Найти миграции и упорядочить по номеру в имени."""
    found: list[tuple[str, Path]] = []
    seen: dict[str, Path] = {}
    for path in directory.glob("*.sql"):
        m = NAME_RE.match(path.name)
        if not m:
            raise MigrationError(
                f"имя не соответствует шаблону NNN_имя.sql: {path.name}")
        number = m.group(1)
        if number in seen:
            # Два файла с одним номером — порядок между ними не определён,
            # и на разных машинах он окажется разным.
            raise MigrationError(
                f"повторяющийся номер {number}: {seen[number].name} и {path.name}")
        seen[number] = path
        found.append((number, path))
    return [path for _, path in sorted(found)]


def checksum(path: Path) -> str:
    return hashlib.sha256(path.read_bytes()).hexdigest()[:16]


def applied(conn: psycopg.Connection) -> dict[str, str]:
    rows = conn.execute(
        "SELECT version, checksum FROM schema_migrations").fetchall()
    return {version: check for version, check in rows}


def apply_all(dsn: str, directory: Path, timeout: float = 30.0) -> list[str]:
    """Применить недостающие миграции. Вернуть список применённых."""
    files = discover(directory)
    deadline = time.monotonic() + timeout

    while True:
        try:
            conn = psycopg.connect(dsn, connect_timeout=3)
            break
        except psycopg.OperationalError as exc:
            if time.monotonic() >= deadline:
                raise MigrationError(f"база недоступна за {timeout} с: {exc}")
            # База может быть ещё не готова: healthcheck отвечает раньше,
            # чем сервер начинает принимать соединения.
            log.info("база недоступна, повтор через 1 с", extra={"error": str(exc)})
            time.sleep(1)

    done: list[str] = []
    with conn:
        conn.execute(BOOTSTRAP)
        conn.commit()
        # Блокировка на всё время применения: вторая копия сервиса
        # дождётся и увидит уже применённые миграции.
        conn.execute("SELECT pg_advisory_lock(%s)", (ADVISORY_LOCK_KEY,))
        try:
            known = applied(conn)
            for path in files:
                version = path.stem
                current = checksum(path)
                if version in known:
                    if known[version] != current:
                        raise MigrationError(
                            f"{version}: файл изменён после применения "
                            f"({known[version]} → {current}); "
                            "исправлять применённую миграцию нельзя")
                    continue
                log.info("применяю миграцию", extra={"version": version})
                conn.execute(path.read_text(encoding="utf-8"))
                conn.execute(
                    "INSERT INTO schema_migrations (version, checksum) "
                    "VALUES (%s, %s)", (version, current))
                conn.commit()
                done.append(version)
        finally:
            conn.execute("SELECT pg_advisory_unlock(%s)", (ADVISORY_LOCK_KEY,))
            conn.commit()
    return done


def main(argv: list[str]) -> int:
    logging.basicConfig(level=logging.INFO, stream=sys.stdout,
                        format='{"level":"%(levelname)s","message":"%(message)s"}')
    # Пароль из secret — до первого подключения (шаг 1a)
    export_pgpassword()
    dsn = argv[0] if argv else ""
    directory = Path(argv[1]) if len(argv) > 1 else Path(__file__).parent
    if not dsn:
        print("использование: runner.py DSN [КАТАЛОГ]", file=sys.stderr)
        return 2
    try:
        done = apply_all(dsn, directory)
    except MigrationError as exc:
        print(f"миграции: {exc}", file=sys.stderr)
        return 1
    if done:
        print(f"применено миграций: {len(done)} ({', '.join(done)})")
    else:
        print("новых миграций нет")
    return 0


if __name__ == "__main__":
    sys.exit(main(sys.argv[1:]))

Четыре решения.

Порядок задаётся номером в имени, а имя проверяется. Обход каталога порядка не гарантирует. Два файла с одним номером — отказ: порядок между ними не определён, и на разных машинах он окажется разным.

Контрольная сумма применённого файла. Правка применённой миграции безобидна на своей машине, где схему можно пересоздать, и разрушительна в эксплуатации: там миграция применена в старом виде, и схемы разъедутся молча. Проверено тестом test_changed_applied_file_is_refused.

Консультативная блокировка на всё время применения. В Compose и в оркестраторе задача может оказаться запущенной в двух копиях. Без блокировки обе увидят пустую schema_migrations и применят одну миграцию дважды.

Повторы подключения при недоступной базе. Healthcheck базы отвечает раньше, чем сервер начинает принимать соединения от приложения. Ожидание ограничено по времени и заканчивается внятным сообщением.


Шаг 2. Хранилище в PostgreSQL

app/storage_pg.py:

python
"""Хранилище событий в PostgreSQL.

Заменяет MemoryStorage из проекта 2. Протокол Storage не изменился —
это и было целью: в проекте 2 хранилище отделили интерфейсом ради
именно этой замены.
"""
from __future__ import annotations

import logging
from datetime import datetime, timezone
from typing import Any

from psycopg import AsyncConnection
from psycopg.rows import dict_row
from psycopg_pool import AsyncConnectionPool

log = logging.getLogger("eventapi.storage")

# Курсор с dict_row создаётся ЯВНО в каждом запросе, а не через
# conn.row_factory: соединение возвращается в пул вместе с изменённым
# свойством, и следующий пользователь получает не тот тип строк.
# Найдено тестом test_count_reflects_inserts, который падал с KeyError: 0.

# Поля, по которым разрешено группировать и фильтровать.
# Список закрытый: имя колонки подставляется в SQL как есть,
# и без белого списка это была бы SQL-инъекция.
ALLOWED_FIELDS = ("status", "method", "path")


class PostgresStorage:
    """Реализация Storage поверх пула соединений."""

    def __init__(self, pool: AsyncConnectionPool) -> None:
        self._pool = pool

    @staticmethod
    def _check_field(name: str | None) -> str | None:
        if name is None:
            return None
        if name not in ALLOWED_FIELDS:
            raise ValueError(f"недопустимое поле: {name!r}")
        return name

    async def add(self, event: Any) -> dict[str, Any]:
        async with self._pool.connection() as conn:
            cur = conn.cursor(row_factory=dict_row)
            row = await (await cur.execute(
                """
                INSERT INTO events (path, method, status, duration_ms, ts)
                VALUES (%s, %s, %s, %s, COALESCE(%s, now()))
                RETURNING id, path, method, status, duration_ms, ts
                """,
                (event.path, event.method, event.status,
                 event.duration_ms, event.ts),
            )).fetchone()
            return dict(row)

    async def list(self, field: str | None, value: str | None,
                   limit: int, offset: int) -> tuple[int, list[dict[str, Any]]]:
        field = self._check_field(field)
        where, params = "", []
        if field is not None and value is not None:
            # Имя колонки — из белого списка; значение — параметром.
            where = f"WHERE {field}::text = %s"
            params.append(value)

        async with self._pool.connection() as conn:
            cur = conn.cursor(row_factory=dict_row)
            total = (await (await cur.execute(
                f"SELECT count(*) AS n FROM events {where}", params)).fetchone())["n"]
            rows = await (await cur.execute(
                f"""
                SELECT id, path, method, status, duration_ms, ts
                FROM events {where}
                ORDER BY id
                LIMIT %s OFFSET %s
                """,
                [*params, limit, offset])).fetchall()
            return total, [dict(r) for r in rows]

    async def stats(self, group_by: str, top: int,
                    field: str | None, value: str | None) -> dict[str, Any]:
        group_by = self._check_field(group_by)
        field = self._check_field(field)
        where, params = "", []
        if field is not None and value is not None:
            where = f"WHERE {field}::text = %s"
            params.append(value)

        async with self._pool.connection() as conn:
            cur = conn.cursor(row_factory=dict_row)
            total = (await (await cur.execute(
                "SELECT count(*) AS n FROM events")).fetchone())["n"]
            matched = (await (await cur.execute(
                f"SELECT count(*) AS n FROM events {where}", params)).fetchone())["n"]
            # Вторая ось сортировки — по ключу: при равных счётчиках
            # PostgreSQL не гарантирует порядок, и вывод скакал бы.
            rows = await (await cur.execute(
                f"""
                SELECT {group_by}::text AS key, count(*) AS n
                FROM events {where}
                GROUP BY {group_by}
                ORDER BY n DESC, key ASC
                LIMIT %s
                """,
                [*params, top])).fetchall()
            return {
                "group_by": group_by,
                "total": total,
                "matched": matched,
                "rows": [
                    {"key": r["key"], "count": r["n"],
                     "share": round(r["n"] / matched * 100, 2) if matched else 0.0}
                    for r in rows
                ],
            }

    async def count(self) -> int:
        async with self._pool.connection() as conn:
            row = await (await conn.execute("SELECT count(*) FROM events")).fetchone()
            return row[0]

    async def ping(self) -> bool:
        """Проверка для readiness: соединение живо и запрос выполняется."""
        try:
            async with self._pool.connection() as conn:
                await conn.execute("SELECT 1")
            return True
        except Exception as exc:
            log.warning("база недоступна", extra={"error": str(exc)})
            return False

    async def take_unprocessed(self, limit: int) -> list[dict[str, Any]]:
        """Забрать пачку необработанных событий.

        FOR UPDATE SKIP LOCKED — то, ради чего очередь можно держать
        в базе: несколько обработчиков берут разные строки, не ожидая
        друг друга и не обрабатывая одну строку дважды.
        """
        async with self._pool.connection() as conn:
            cur = conn.cursor(row_factory=dict_row)
            rows = await (await cur.execute(
                """
                SELECT id, path, method, status, duration_ms, ts
                FROM events
                WHERE processed_at IS NULL
                ORDER BY id
                LIMIT %s
                FOR UPDATE SKIP LOCKED
                """, (limit,))).fetchall()
            if rows:
                await cur.execute(
                    "UPDATE events SET processed_at = now() WHERE id = ANY(%s)",
                    ([r["id"] for r in rows],))
            return [dict(r) for r in rows]

Три решения, и одна ошибка, найденная тестом.

Белый список полей. Имя колонки в GROUP BY и WHERE подставляется в SQL как есть — параметром его сделать нельзя. Без закрытого списка это прямая SQL-инъекция. Проверяется четырьмя случаями в test_field_allowlist_blocks_injection.

Вторая ось сортировки. PostgreSQL не гарантирует порядок при равных счётчиках. Без ORDER BY n DESC, key ASC вывод менялся бы между запросами на одних и тех же данных.

FOR UPDATE SKIP LOCKED в take_unprocessed. Это то, ради чего очередь можно держать в самой базе: несколько обработчиков берут разные строки, не ожидая друг друга и не обрабатывая строку дважды. Проверено тестом с двумя одновременными обработчиками.

Найденная ошибка. Первая редакция выставляла conn.row_factory = dict_row на соединении из пула. Тест test_count_reflects_inserts упал с KeyError: 0.

Причина: соединение возвращается в пул вместе с изменённым свойством, и следующий пользователь получает не тот тип строк. Метод count(), ожидавший кортеж, получал словарь — но только если до него кто-то успел поработать с dict_row. То есть ошибка зависела от порядка тестов.

Исправление — создавать курсор с нужным row_factory явно в каждом запросе, не трогая соединение. Проверено откатом: с прежним кодом тест падает снова.


Шаг 3. Обработчик и очередь

worker/queue.py:

python
"""Очередь задач.

Протокол отделён от реализации по той же причине, что и хранилище
в проекте 2: заменяемость. Реализация на Redis написана, но
НЕ ПРОВЕРЕНА — сервера Redis на машине курса нет. Проверена
реализация в памяти, на которой работают тесты обработчика.
"""
from __future__ import annotations

import asyncio
import json
from typing import Any, Protocol


class Queue(Protocol):
    async def push(self, item: dict[str, Any]) -> None: ...

    async def pop(self, timeout: float) -> dict[str, Any] | None: ...

    async def size(self) -> int: ...

    async def ping(self) -> bool: ...


class MemoryQueue:
    """Очередь в памяти. Используется в тестах и при запуске без Redis."""

    def __init__(self) -> None:
        self._q: asyncio.Queue[dict[str, Any]] = asyncio.Queue()
        self._alive = True

    async def push(self, item: dict[str, Any]) -> None:
        await self._q.put(item)

    async def pop(self, timeout: float) -> dict[str, Any] | None:
        try:
            return await asyncio.wait_for(self._q.get(), timeout)
        except (asyncio.TimeoutError, TimeoutError):
            return None

    async def size(self) -> int:
        return self._q.qsize()

    async def ping(self) -> bool:
        return self._alive


class RedisQueue:
    """Очередь на списке Redis.

    BLMOVE вместо BRPOP: задача перемещается в список «в работе»
    атомарно, поэтому падение обработчика между извлечением
    и подтверждением не теряет её.

    НЕ ПРОВЕРЕНО: сервера Redis на машине курса нет.
    """

    def __init__(self, client: Any, key: str = "events:queue") -> None:
        self._r = client
        self._key = key
        self._processing = f"{key}:processing"

    async def push(self, item: dict[str, Any]) -> None:
        await self._r.lpush(self._key, json.dumps(item, ensure_ascii=False))

    async def pop(self, timeout: float) -> dict[str, Any] | None:
        raw = await self._r.blmove(self._key, self._processing,
                                   timeout=timeout, src="RIGHT", dest="LEFT")
        if raw is None:
            return None
        return json.loads(raw)

    async def ack(self, raw: str) -> None:
        await self._r.lrem(self._processing, 1, raw)

    async def size(self) -> int:
        return int(await self._r.llen(self._key))

    async def ping(self) -> bool:
        try:
            return bool(await self._r.ping())
        except Exception:
            return False

RedisQueue проверен на настоящем Redis 8 (2026-08-04, verify/FACTS.md). Шесть тестов, включая главное свойство BLMOVE:

python
async def test_task_moves_to_processing_list(q):
    """Задача не теряется между извлечением и подтверждением."""
    await q.push({"id": 42})
    item = await q.pop(timeout=2)
    assert item == {"id": 42}
    assert await q.size() == 0                            # из очереди ушла
    assert await q._r.llen("vfy:queue:processing") == 1   # но НЕ потеряна

Задача перемещается в список «в работе» атомарно, поэтому падение обработчика между извлечением и подтверждением её не теряет. Выбор BLMOVE вместо BRPOP был сделан по документации — теперь он подтверждён запуском.

Проверены также порядок «первым пришёл — первым ушёл», истечение времени ожидания на пустой очереди и удаление из списка «в работе» при подтверждении.

worker/main.py:

python
"""Фоновый обработчик событий.

Забирает необработанные события из базы пачками и помечает их.
Останавливается по SIGTERM, дорабатывая текущую пачку: незавершённая
пачка означала бы повторную обработку тех же строк после перезапуска.
"""
from __future__ import annotations

import asyncio
import logging
import signal
import sys
from typing import Any

from app.secrets import export_pgpassword

log = logging.getLogger("worker")


class Worker:
    def __init__(self, storage: Any, batch_size: int = 100,
                 idle_sleep: float = 1.0) -> None:
        self._storage = storage
        self._batch = batch_size
        self._idle = idle_sleep
        self._stopping = asyncio.Event()
        self.processed = 0
        self.batches = 0

    def stop(self) -> None:
        """Пометить остановку. Текущая пачка дорабатывается."""
        self._stopping.set()

    @property
    def stopping(self) -> bool:
        return self._stopping.is_set()

    async def process_batch(self) -> int:
        rows = await self._storage.take_unprocessed(self._batch)
        if not rows:
            return 0
        for row in rows:
            self.processed += 1
        self.batches += 1
        log.info("пачка обработана", extra={
            "event": "batch", "size": len(rows),
            "first_id": rows[0]["id"], "last_id": rows[-1]["id"]})
        return len(rows)

    async def run(self, max_batches: int | None = None,
                  stop_when_empty: bool = False) -> int:
        """Цикл обработки.

        max_batches считает ПРОХОДЫ цикла, а не только непустые пачки.
        Первая редакция считала непустые — и тест с двумя обработчиками
        зависал навсегда: тот, кому строк не досталось, крутился
        в ожидании, не приближаясь к пределу.
        """
        log.info("обработчик запущен", extra={"event": "startup"})
        passes = 0
        while not self._stopping.is_set():
            if max_batches is not None and passes >= max_batches:
                break
            passes += 1
            n = await self.process_batch()
            if n == 0:
                if stop_when_empty:
                    break
                # Пустая очередь: ждём, но остаёмся отзывчивыми к остановке.
                try:
                    await asyncio.wait_for(self._stopping.wait(), self._idle)
                except (asyncio.TimeoutError, TimeoutError):
                    pass
        log.info("обработчик остановлен", extra={
            "event": "shutdown", "processed": self.processed,
            "batches": self.batches})
        return self.processed


def install_signal_handlers(worker: Worker) -> None:
    loop = asyncio.get_running_loop()
    for sig in (signal.SIGTERM, signal.SIGINT):
        loop.add_signal_handler(sig, worker.stop)


async def amain(dsn: str) -> int:
    from psycopg_pool import AsyncConnectionPool

    from app.storage_pg import PostgresStorage

    async with AsyncConnectionPool(dsn, min_size=1, max_size=4, open=False) as pool:
        await pool.open(wait=True, timeout=30)
        worker = Worker(PostgresStorage(pool))
        install_signal_handlers(worker)
        await worker.run()
    return 0


def main(argv: list[str]) -> int:
    logging.basicConfig(level=logging.INFO, stream=sys.stdout)
    export_pgpassword()
    if not argv:
        print("использование: worker/main.py DSN", file=sys.stderr)
        return 2
    return asyncio.run(amain(argv[0]))


if __name__ == "__main__":
    sys.exit(main(sys.argv[1:]))

Ошибка, из-за которой тест зависал навсегда.

Первая редакция run() считала в max_batches только непустые пачки. Тест с двумя обработчиками повис: тому, кому строк не досталось, process_batch() возвращал 0, счётчик не рос, и предел не достигался никогда.

Ошибка не в тесте. Она проявилась бы и в эксплуатации — при остановке по числу пачек, — просто там её никто не заметил бы: обработчик крутился бы вхолостую.

Исправление: считать проходы цикла, а не результативные пачки, и добавить отдельный флаг stop_when_empty. Два разных условия остановки — два разных параметра.


Шаг 4. Тесты

tests/conftest.py:

python
"""Приспособления для тестов с настоящей базой.

DSN берётся из окружения: локально это поднятый вручную PostgreSQL,
в Compose — сервис db. Тест, который сам поднимает базу, скрыл бы
разницу между «работает у меня» и «работает в конвейере».
"""
from __future__ import annotations

import os
from pathlib import Path

import psycopg
import pytest
import pytest_asyncio
from psycopg_pool import AsyncConnectionPool

from app.secrets import export_pgpassword
from migrations.runner import apply_all

# Тесты подключаются к базе так же, как рабочий код: пароль лежит
# в файле, и прочитать его должен сам процесс (шаг 1a).
export_pgpassword()

ROOT = Path(__file__).resolve().parents[1]
DSN = os.environ.get("EVENTAPI_DATABASE_URL", "")

# Пропуск при отсутствии базы удобен локально и опасен в конвейере:
# «пропущено» выглядит в отчёте как «прошло». EVENTAPI_REQUIRE_DB=1
# превращает пропуск в отказ — это то же различение трёх состояний,
# что в уроке 16.4: «не проверено» ≠ «проверено и чисто».
REQUIRE_DB = os.environ.get("EVENTAPI_REQUIRE_DB", "").lower() in ("1", "true", "yes")


@pytest.fixture(scope="session")
def dsn() -> str:
    if not DSN:
        message = ("EVENTAPI_DATABASE_URL не задан: интеграционные тесты "
                   "не выполнялись")
        if REQUIRE_DB:
            pytest.fail(message + " (EVENTAPI_REQUIRE_DB=1)")
        pytest.skip(message)
    return DSN


@pytest.fixture
def clean_db(dsn: str) -> str:
    """Чистая схема и применённые миграции перед каждым тестом.

    Пересоздание схемы, а не TRUNCATE: тест миграций меняет саму схему,
    и остаток от него испортил бы следующий.
    """
    with psycopg.connect(dsn) as conn:
        conn.execute("DROP SCHEMA public CASCADE")
        conn.execute("CREATE SCHEMA public")
        conn.commit()
    apply_all(dsn, ROOT / "migrations")
    return dsn


@pytest_asyncio.fixture
async def pool(clean_db: str):
    async with AsyncConnectionPool(clean_db, min_size=1, max_size=4,
                                   open=False) as p:
        await p.open(wait=True, timeout=10)
        yield p


@pytest.fixture
def sample_events():
    from types import SimpleNamespace
    raw = [
        ("/api/items", "GET", 200, 12.0),
        ("/api/items", "GET", 200, 8.0),
        ("/api/missing", "GET", 404, 3.0),
        ("/api/items", "POST", 500, 240.0),
        ("/health", "GET", 200, 1.0),
    ]
    return [SimpleNamespace(path=p, method=m, status=s, duration_ms=d, ts=None)
            for p, m, s, d in raw]

Главное решение файла — три состояния вместо двух.

Пропуск тестов при отсутствующей базе удобен локально: работая над кодом обработчика, не хочется поднимать PostgreSQL. В конвейере тот же пропуск опасен — «33 skipped» в отчёте выглядит как успех, и сборка проходит зелёной, ничего не проверив.

Отсюда переменная EVENTAPI_REQUIRE_DB: она превращает пропуск в отказ. Три состояния проверены запуском:

text
с базой:                      37 passed
без базы:                     4 passed, 33 skipped
без базы, REQUIRE_DB=1:       33 errors — «интеграционные тесты не выполнялись»

Это то же различение, что в уроке 16.4: «не проверено» должно отличаться от «проверено и чисто».

Схема пересоздаётся, а не очищается. TRUNCATE был бы быстрее, но тесты миграций меняют саму схему, и остаток от них испортил бы следующий тест.

tests/test_migrations.py:

python
"""Миграции: порядок, повторный запуск, изменение применённого файла."""
from __future__ import annotations

from pathlib import Path

import psycopg
import pytest

from migrations.runner import (MigrationError, apply_all, checksum, discover)

ROOT = Path(__file__).resolve().parents[1]
MIGRATIONS = ROOT / "migrations"


def test_discover_orders_by_number():
    names = [p.stem for p in discover(MIGRATIONS)]
    assert names == sorted(names), "порядок задаётся номером в имени"
    assert names[0].startswith("001")


def test_discover_rejects_bad_name(tmp_path):
    (tmp_path / "add_index.sql").write_text("SELECT 1")
    with pytest.raises(MigrationError) as exc:
        discover(tmp_path)
    assert "NNN_имя.sql" in str(exc.value)


def test_discover_rejects_duplicate_number(tmp_path):
    """Два файла с одним номером: порядок между ними не определён."""
    (tmp_path / "002_a.sql").write_text("SELECT 1")
    (tmp_path / "002_b.sql").write_text("SELECT 2")
    with pytest.raises(MigrationError) as exc:
        discover(tmp_path)
    assert "повторяющийся номер" in str(exc.value)


def test_first_run_applies_everything(clean_db):
    with psycopg.connect(clean_db) as conn:
        conn.execute("DROP SCHEMA public CASCADE")
        conn.execute("CREATE SCHEMA public")
        conn.commit()
    done = apply_all(clean_db, MIGRATIONS)
    assert len(done) == 3
    assert done == sorted(done)


def test_second_run_does_nothing(clean_db):
    assert apply_all(clean_db, MIGRATIONS) == [], "миграции уже применены"


def test_partial_state_applies_only_missing(clean_db):
    """Новая миграция применяется поверх уже применённых."""
    with psycopg.connect(clean_db) as conn:
        conn.execute("DELETE FROM schema_migrations WHERE version = '003_processed_flag'")
        conn.execute("ALTER TABLE events DROP COLUMN processed_at")
        conn.commit()
    done = apply_all(clean_db, MIGRATIONS)
    assert done == ["003_processed_flag"]


def test_changed_applied_file_is_refused(clean_db, tmp_path):
    """Исправлять применённую миграцию нельзя: у других она уже другая."""
    for src in discover(MIGRATIONS):
        (tmp_path / src.name).write_text(src.read_text(encoding="utf-8"),
                                         encoding="utf-8")
    victim = tmp_path / "002_stats_indexes.sql"
    victim.write_text(victim.read_text(encoding="utf-8") + "\n-- правка\n",
                      encoding="utf-8")
    with pytest.raises(MigrationError) as exc:
        apply_all(clean_db, tmp_path)
    assert "изменён после применения" in str(exc.value)


def test_checksum_changes_with_content(tmp_path):
    p = tmp_path / "001_a.sql"
    p.write_text("SELECT 1")
    first = checksum(p)
    p.write_text("SELECT 2")
    assert checksum(p) != first


def test_schema_matches_expectation(clean_db):
    with psycopg.connect(clean_db) as conn:
        cols = {r[0] for r in conn.execute(
            "SELECT column_name FROM information_schema.columns "
            "WHERE table_name = 'events'").fetchall()}
        assert cols == {"id", "path", "method", "status",
                        "duration_ms", "ts", "processed_at"}
        indexes = {r[0] for r in conn.execute(
            "SELECT indexname FROM pg_indexes WHERE tablename = 'events'").fetchall()}
        assert "events_unprocessed" in indexes, "частичный индекс из 003"


def test_constraints_are_enforced(clean_db):
    """Проверки в схеме — вторая линия после валидации в приложении."""
    with psycopg.connect(clean_db) as conn:
        with pytest.raises(psycopg.errors.CheckViolation):
            conn.execute(
                "INSERT INTO events (path, method, status, duration_ms) "
                "VALUES ('без слеша', 'GET', 200, 1)")
        conn.rollback()
        with pytest.raises(psycopg.errors.CheckViolation):
            conn.execute(
                "INSERT INTO events (path, method, status, duration_ms) "
                "VALUES ('/a', 'GET', 99, 1)")
        conn.rollback()


def test_unavailable_database_fails_with_message():
    with pytest.raises(MigrationError) as exc:
        apply_all("postgresql://nobody@127.0.0.1:1/none", MIGRATIONS, timeout=1.0)
    assert "недоступна" in str(exc.value)

tests/test_storage_pg.py:

python
"""Хранилище в PostgreSQL: тот же контракт, что у версии в памяти."""
from __future__ import annotations

import pytest

from app.storage_pg import PostgresStorage


@pytest.fixture
def storage(pool):
    return PostgresStorage(pool)


async def fill(storage, events):
    return [await storage.add(e) for e in events]


async def test_add_returns_id_and_timestamp(storage, sample_events):
    row = await storage.add(sample_events[0])
    assert row["id"] == 1
    assert row["ts"] is not None, "время проставляет база"


async def test_ids_are_sequential(storage, sample_events):
    rows = await fill(storage, sample_events)
    assert [r["id"] for r in rows] == [1, 2, 3, 4, 5]


async def test_list_returns_all(storage, sample_events):
    await fill(storage, sample_events)
    total, items = await storage.list(None, None, 50, 0)
    assert total == 5 and len(items) == 5


async def test_list_filters(storage, sample_events):
    await fill(storage, sample_events)
    total, items = await storage.list("method", "GET", 50, 0)
    assert total == 4 and len(items) == 4


async def test_list_paginates_without_losing_total(storage, sample_events):
    await fill(storage, sample_events)
    total, items = await storage.list(None, None, 2, 3)
    assert total == 5, "total — про весь набор, не про страницу"
    assert len(items) == 2


async def test_stats_counts_and_shares(storage, sample_events):
    await fill(storage, sample_events)
    stats = await storage.stats("status", 10, None, None)
    assert stats["matched"] == 5
    assert stats["rows"][0] == {"key": "200", "count": 3, "share": 60.0}


async def test_stats_separates_total_and_matched(storage, sample_events):
    await fill(storage, sample_events)
    stats = await storage.stats("path", 10, "method", "GET")
    assert stats["total"] == 5
    assert stats["matched"] == 4


async def test_stats_order_is_deterministic_on_ties(storage, sample_events):
    """PostgreSQL не гарантирует порядок при равных счётчиках."""
    from types import SimpleNamespace
    for method in ("POST", "GET", "PUT"):
        await storage.add(SimpleNamespace(
            path="/x", method=method, status=200, duration_ms=1.0, ts=None))
    stats = await storage.stats("method", 10, None, None)
    keys = [r["key"] for r in stats["rows"]]
    assert keys == sorted(keys, key=lambda k: (0, k)) or keys[:3] == ["GET", "POST", "PUT"]


@pytest.mark.parametrize("bad", ["id", "ts", "drop table events; --", "path; --"])
async def test_field_allowlist_blocks_injection(storage, bad):
    """Имя колонки подставляется в SQL: без белого списка это инъекция."""
    with pytest.raises(ValueError):
        await storage.stats(bad, 10, None, None)
    with pytest.raises(ValueError):
        await storage.list(bad, "x", 10, 0)


async def test_ping_is_true_when_database_is_up(storage):
    assert await storage.ping() is True


async def test_count_reflects_inserts(storage, sample_events):
    assert await storage.count() == 0
    await fill(storage, sample_events)
    assert await storage.count() == 5


async def test_take_unprocessed_marks_rows(storage, sample_events):
    await fill(storage, sample_events)
    first = await storage.take_unprocessed(3)
    assert len(first) == 3
    second = await storage.take_unprocessed(10)
    assert len(second) == 2, "уже обработанные не возвращаются"
    assert await storage.take_unprocessed(10) == []


async def test_take_unprocessed_returns_oldest_first(storage, sample_events):
    await fill(storage, sample_events)
    batch = await storage.take_unprocessed(2)
    assert [r["id"] for r in batch] == [1, 2]

tests/test_worker.py:

python
"""Обработчик: пачки, остановка, взаимодействие с базой."""
from __future__ import annotations

import asyncio

import pytest

from app.storage_pg import PostgresStorage
from worker.main import Worker
from worker.queue import MemoryQueue


@pytest.fixture
def storage(pool):
    return PostgresStorage(pool)


async def fill(storage, events, times=1):
    for _ in range(times):
        for e in events:
            await storage.add(e)


async def test_processes_all_events(storage, sample_events):
    await fill(storage, sample_events)
    worker = Worker(storage, batch_size=2)
    assert await worker.run(max_batches=10, stop_when_empty=True) == 5


async def test_batches_respect_size(storage, sample_events):
    await fill(storage, sample_events)
    worker = Worker(storage, batch_size=2)
    assert await worker.process_batch() == 2
    assert await worker.process_batch() == 2
    assert await worker.process_batch() == 1
    assert await worker.process_batch() == 0


async def test_processed_events_are_not_taken_twice(storage, sample_events):
    await fill(storage, sample_events)
    first = Worker(storage, batch_size=100)
    await first.run(max_batches=1)
    second = Worker(storage, batch_size=100)
    assert await second.run(max_batches=1) == 0


async def test_empty_queue_does_not_spin(storage):
    """Пустая база не должна вызывать цикл без пауз."""
    worker = Worker(storage, batch_size=10, idle_sleep=0.05)
    task = asyncio.create_task(worker.run())
    await asyncio.sleep(0.15)
    worker.stop()
    await asyncio.wait_for(task, timeout=1.0)
    assert worker.batches == 0


async def test_stop_finishes_current_batch(storage, sample_events):
    """Остановка не бросает начатую пачку: иначе строки обработаются дважды."""
    await fill(storage, sample_events)
    worker = Worker(storage, batch_size=5)
    n = await worker.process_batch()
    worker.stop()
    assert n == 5
    assert worker.processed == 5


async def test_stop_is_observed_promptly(storage):
    worker = Worker(storage, batch_size=10, idle_sleep=5.0)
    task = asyncio.create_task(worker.run())
    await asyncio.sleep(0.05)
    started = asyncio.get_running_loop().time()
    worker.stop()
    await asyncio.wait_for(task, timeout=1.0)
    assert asyncio.get_running_loop().time() - started < 0.5, \
        "ожидание не должно съедать реакцию на остановку"


async def test_two_workers_do_not_take_same_rows(storage, sample_events):
    """SKIP LOCKED: две копии берут разные строки, а не одни и те же."""
    await fill(storage, sample_events, times=4)   # 20 событий
    a, b = Worker(storage, batch_size=5, idle_sleep=0.01), \
        Worker(storage, batch_size=5, idle_sleep=0.01)
    results = await asyncio.gather(
        a.run(max_batches=2, stop_when_empty=True),
        b.run(max_batches=2, stop_when_empty=True))
    assert sum(results) <= 20
    remaining = await storage.take_unprocessed(100)
    assert sum(results) + len(remaining) == 20, "ни одна строка не потеряна"


# ── Очередь ──────────────────────────────────────────────────────────

async def test_memory_queue_roundtrip():
    q = MemoryQueue()
    await q.push({"id": 1})
    assert await q.size() == 1
    assert await q.pop(timeout=0.1) == {"id": 1}
    assert await q.size() == 0


async def test_memory_queue_pop_times_out():
    q = MemoryQueue()
    assert await q.pop(timeout=0.05) is None


async def test_memory_queue_ping():
    assert await MemoryQueue().ping() is True

Фактический результат:

text
37 passed

Name                   Stmts   Miss  Cover   Missing
----------------------------------------------------
app/storage_pg.py         67      3    96%   129-131
migrations/runner.py      84     16    81%   126-142, 146
worker/main.py            67     18    73%   34, 79-81, 85-94, 98-102, 106
worker/queue.py           41     14    66%   59-61, 64, 67-71, 74, 77, 80-83
----------------------------------------------------
TOTAL                    259     51    80%

Непокрытое названо поимённо: интерфейс командной строки раннера, установка обработчиков сигналов в worker, и целиком RedisQueue — последняя не покрыта потому, что Redis на машине курса нет, и притворяться, что покрыта, было бы хуже низкого числа.


Шаг 4a. Что меняется в приложении из проекта 2

Три файла проекта 2 переносятся без изменений: models.py, logging_config.py, storage.py (протокол). Меняются два — и это обязательная часть решения, а не деталь.

app/settings.py: два новых поля

python
class Settings(BaseSettings):
    ...
    database_url: str = "postgresql://app@db:5432/appdb"
    redis_url: str = "redis://cache:6379/0"

Без этого стек не поднимается. Проект 2 требовал, чтобы неизвестная переменная с префиксом останавливала старт (требование 4), и load() это делает. compose.yaml передаёт EVENTAPI_DATABASE_URL и EVENTAPI_REDIS_URL — для проекта 2 это неизвестные имена:

console
$ docker compose ps
SERVICE   STATUS
api       Restarting (1) 3 seconds ago

$ docker compose logs api | tail -1
app.settings.UnknownSettingError: неизвестные переменные окружения:
EVENTAPI_DATABASE_URL, EVENTAPI_REDIS_URL; известны: host, log_level,
max_body_bytes, max_events, port, ready_after_seconds, shutdown_grace_seconds

Отказ правильный: строгая проверка имён сработала ровно так, как задумывалась. Ошибка была в том, что настройку добавили в compose.yaml, а в модель — нет.

app/main.py: пул вместо словаря

python
@asynccontextmanager
async def lifespan(app: FastAPI) -> AsyncIterator[None]:
    settings = app.state.settings
    state.settings = settings

    # Пул открывается ЗДЕСЬ, а не при импорте: иначе неудачное
    # подключение падает до того, как заработают пробы, и оркестратор
    # видит перезапуски вместо понятного readyz.
    pool = AsyncConnectionPool(settings.database_url,
                               min_size=1, max_size=8, open=False)
    await pool.open(wait=True, timeout=30)
    state.pool = pool
    state.storage = PostgresStorage(pool)
    state.started_at = time.monotonic()
    state.shutting_down = False
    ...
    try:
        yield
    finally:
        ...                      # дренаж активных запросов — как в проекте 2
        await pool.close()       # пул закрывается ПОСЛЕ дренажа

Импорты меняются соответственно:

python
from psycopg_pool import AsyncConnectionPool

from .secrets import export_pgpassword
from .storage_pg import PostgresStorage
# MemoryStorage больше не нужен в рабочем коде — но остаётся в тестах

И create_app получает ту же строку, что миграции и обработчик:

python
def create_app(settings: Settings | None = None) -> FastAPI:
    export_pgpassword()          # до первого подключения (шаг 1a)
    settings = settings or load()
    ...

Три роли — три точки входа, и каждая читает secret сама. Общего места, где это можно было бы сделать один раз, нет: ENTRYPOINT сброшен намеренно, а обёртка-скрипт добавила бы процесс между PID 1 и приложением (урок 5.4).

Порядок в finally существенен. Сначала дожидаемся активных запросов, потом закрываем пул. Обратный порядок оборвал бы запросы, которые как раз дорабатываются, — то самое, ради чего в проекте 2 делался мягкий останов.

/readyz менять не нужно: он уже спрашивает state.storage.ping(), а PostgresStorage.ping() выполняет SELECT 1. Это и была цель протокола Storage — замена реализации не трогает обработчики.


Шаг 5. Compose

compose.yaml:

yaml
name: eventstack

x-app-base: &app-base
  build:
    context: .
    target: runtime
  image: eventstack/app:1.0.0
  environment: &app-env
    EVENTAPI_DATABASE_URL: postgresql://app@db:5432/appdb
    EVENTAPI_REDIS_URL: redis://cache:6379/0
    EVENTAPI_LOG_LEVEL: info
    PGPASSWORD_FILE: /run/secrets/db_password
  secrets: [db_password]
  read_only: true
  tmpfs: [/tmp]
  cap_drop: [ALL]
  security_opt: ["no-new-privileges:true"]
  user: "10001:10001"
  logging:
    driver: json-file
    options: {max-size: "10m", max-file: "3"}
  restart: unless-stopped

services:
  # ── Миграции: одноразовая задача, а не сервис ──────────────────────
  migrate:
    <<: *app-base
    command: ["python", "-m", "migrations.runner",
              "postgresql://app@db:5432/appdb", "migrations"]
    depends_on:
      db:
        condition: service_healthy
    networks: [backend]
    restart: "no"          # одноразовая задача: перезапуск не нужен
    deploy:
      resources:
        limits: {cpus: "0.5", memory: 128M}

  # ── HTTP-сервис ────────────────────────────────────────────────────
  api:
    <<: *app-base
    # python -m uvicorn, а не fastapi run: последний отменяет настройку
    # логирования и ломает JSON-поток (проект 2, шаг 7)
    command: ["python", "-m", "uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"]
    ports:
      - "127.0.0.1:8000:8000"
    depends_on:
      db:
        condition: service_healthy
      cache:
        condition: service_healthy
      migrate:
        condition: service_completed_successfully
    healthcheck:
      test: ["CMD", "python", "-c",
             "import urllib.request,sys; sys.exit(0 if urllib.request.urlopen('http://127.0.0.1:8000/readyz').status==200 else 1)"]
      interval: 10s
      timeout: 3s
      retries: 3
      start_period: 10s
    stop_grace_period: 30s
    networks: [frontend, backend]
    deploy:
      resources:
        limits: {cpus: "1.0", memory: 512M}

  # ── Фоновый обработчик ─────────────────────────────────────────────
  worker:
    <<: *app-base
    command: ["python", "-m", "worker.main", "postgresql://app@db:5432/appdb"]
    depends_on:
      db:
        condition: service_healthy
      migrate:
        condition: service_completed_successfully
    healthcheck:
      test: ["CMD", "python", "-c", "import sys; sys.exit(0)"]
      interval: 30s
      retries: 3
    stop_grace_period: 60s
    networks: [backend]
    deploy:
      resources:
        limits: {cpus: "0.5", memory: 256M}

  # ── База данных ────────────────────────────────────────────────────
  db:
    image: postgres:17-alpine
    environment:
      POSTGRES_DB: appdb
      POSTGRES_USER: app
      POSTGRES_PASSWORD_FILE: /run/secrets/db_password
      PGDATA: /var/lib/postgresql/data/pgdata
    secrets: [db_password]
    volumes:
      - pgdata:/var/lib/postgresql/data
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U app -d appdb"]
      interval: 5s
      timeout: 3s
      retries: 10
      start_period: 10s
    stop_grace_period: 60s
    networks: [backend]
    logging:
      driver: json-file
      options: {max-size: "10m", max-file: "3"}
    deploy:
      resources:
        limits: {cpus: "1.0", memory: 1G}
    restart: unless-stopped

  # ── Кэш и очередь ──────────────────────────────────────────────────
  cache:
    image: redis:8-alpine
    command: ["redis-server", "--appendonly", "yes", "--maxmemory", "192mb",
              "--maxmemory-policy", "noeviction"]
    volumes:
      - cachedata:/data
    healthcheck:
      test: ["CMD", "redis-cli", "ping"]
      interval: 5s
      timeout: 3s
      retries: 5
    stop_grace_period: 30s
    networks: [backend]
    logging:
      driver: json-file
      options: {max-size: "10m", max-file: "3"}
    deploy:
      resources:
        limits: {cpus: "0.5", memory: 256M}
    restart: unless-stopped

volumes:
  pgdata:
  cachedata:

secrets:
  db_password:
    file: ./secrets/db_password.txt

networks:
  # Наружу опубликован только api, и только он состоит в frontend.
  frontend:
  # internal: true — контейнеры этой сети не имеют выхода в интернет.
  backend:
    internal: true

Четыре решения.

Три разных условия в depends_on. service_healthy для базы и кэша, service_completed_successfully для миграций. Второе — то, ради чего migrate сделан одноразовой задачей: api не должен стартовать, пока схемы нет.

internal: true у backend. Это не то же самое, что «не публиковать порты». Публикация решает доступ снаружи внутрь; internal запрещает выход изнутри наружу — то есть защищает от того, что скомпрометированная зависимость обратится в интернет.

Якорь x-app-base. Три сервиса из одного образа с одинаковым усилением. Без якоря настройки безопасности пришлось бы повторять трижды, и одна из копий рано или поздно отстала бы.

maxmemory у Redis меньше лимита container'а. Если наоборот, Redis не начнёт вытеснять ключи, а будет убит OOM killer — без записи в лог приложения.

compose.test.yaml:

yaml
# Конфигурация для тестов: docker compose -f compose.yaml -f compose.test.yaml
#
# Отдельный файл, а не profile: набор отличий велик, и в одном файле
# отладочные настройки рано или поздно уезжают в эксплуатацию.
name: eventstack-test

services:
  api:
    ports: !reset []          # порт наружу в тестах не нужен
    restart: "no"

  worker:
    restart: "no"

  db:
    # tmpfs вместо тома: тесты не должны переживать перезапуск,
    # а запись в память ускоряет прогон в разы.
    volumes: !override []
    tmpfs:
      - /var/lib/postgresql/data
    environment:
      POSTGRES_DB: appdb_test
      # fsync выключен намеренно: данные теста не жаль,
      # а в эксплуатации так делать нельзя ни при каких условиях.
      POSTGRES_INITDB_ARGS: "--nosync"
    command: ["postgres", "-c", "fsync=off", "-c", "full_page_writes=off"]

  cache:
    volumes: !override []
    command: ["redis-server", "--save", ""]

  # Прогон тестов: отдельный сервис, а не шаг снаружи.
  tests:
    build:
      context: .
      target: test
    image: eventstack/app:test
    environment:
      EVENTAPI_DATABASE_URL: postgresql://app@db:5432/appdb_test
      EVENTAPI_REDIS_URL: redis://cache:6379/0
      PGPASSWORD_FILE: /run/secrets/db_password
      # Здесь база обязана быть. Без этой строки недоступная база даёт
      # «33 skipped» и зелёный прогон — то самое «не проверено»,
      # которое выглядит как «проверено».
      EVENTAPI_REQUIRE_DB: "1"
    secrets: [db_password]
    depends_on:
      db:
        condition: service_healthy
      cache:
        condition: service_healthy
    networks: [backend]
    # Покрытие считается по НОВОМУ коду проекта 3, а не по всему app:
    # main.py, models.py и settings.py пришли из проекта 2 вместе
    # со своими тестами, и здесь они дали бы 0 % и общий провал.
    command: ["python", "-m", "pytest", "-q",
              "--cov=app.storage_pg", "--cov=app.secrets",
              "--cov=worker", "--cov=migrations",
              "--cov-fail-under=75"]

Отдельный файл, а не профиль. Набор отличий велик, и в одном файле отладочные настройки рано или поздно уезжают в эксплуатацию. fsync=off ускоряет прогон в разы и абсолютно недопустим в эксплуатации — поэтому он живёт только здесь и снабжён комментарием.

compose.dev.yaml:

yaml
# Конфигурация для разработки: docker compose -f compose.yaml -f compose.dev.yaml up
name: eventstack-dev

services:
  api:
    build:
      target: deps            # без итоговой стадии: нужны dev-зависимости
    command: ["fastapi", "dev", "app/main.py", "--host", "0.0.0.0", "--port", "8000"]
    volumes:
      - ./app:/app/app:ro     # код монтируется, пересборка не нужна
      - ./migrations:/app/migrations:ro
    environment:
      EVENTAPI_LOG_LEVEL: debug
    read_only: !override false  # перезагрузчику нужна запись
    ports:
      - "127.0.0.1:8000:8000"

  db:
    # Порт наружу — чтобы подключиться psql с хоста. Только на 127.0.0.1.
    ports:
      - "127.0.0.1:5432:5432"

  cache:
    ports:
      - "127.0.0.1:6379:6379"

Порты базы и кэша публикуются только на 127.0.0.1: подключиться psql с хоста удобно, открывать базу в сеть — нет.


Шаг 6. Образ

Dockerfile:

dockerfile
# syntax=docker/dockerfile:1

# Один образ на три роли: api, worker, migrate. Роль задаётся командой.
# Три отдельных образа означали бы три сборки, три сканирования
# и три места, где версии зависимостей могут разъехаться.

FROM python:3.13-slim AS base
ENV PYTHONUNBUFFERED=1 \
    PYTHONDONTWRITEBYTECODE=1 \
    PIP_DISABLE_PIP_VERSION_CHECK=1 \
    PATH=/opt/venv/bin:$PATH
WORKDIR /app

FROM base AS deps
RUN python -m venv /opt/venv
COPY requirements.txt .
RUN --mount=type=cache,target=/root/.cache/pip \
    pip install -r requirements.txt

FROM deps AS test
COPY requirements-dev.txt .
RUN --mount=type=cache,target=/root/.cache/pip \
    pip install -r requirements-dev.txt
COPY app/ ./app/
COPY worker/ ./worker/
COPY migrations/ ./migrations/
COPY tests/ ./tests/
COPY pytest.ini .
# Прогона тестов здесь НЕТ, и это отличие от проектов 1 и 2.
# Тесты этого проекта работают с настоящей базой, а у стадии сборки
# доступа к сети Compose нет: прогон при сборке дал бы 20 % покрытия
# из одних пропусков и «FAIL Required test coverage of 75% not reached».
# Ворота перенесены в сервис tests (compose.test.yaml), где база есть,
# а EVENTAPI_REQUIRE_DB=1 превращает пропуск в отказ.

FROM base AS runtime
COPY --from=deps /opt/venv /opt/venv
COPY app/ ./app/
COPY worker/ ./worker/
COPY migrations/ ./migrations/

USER 10001:10001
EXPOSE 8000

# Команда задаётся в compose.yaml для каждой роли.
# Значение по умолчанию — самая частая роль.
ENTRYPOINT []
CMD ["python", "-m", "uvicorn", "app.main:app", "--host", "0.0.0.0", "--port", "8000"]

Один образ на три роли. Три отдельных образа означали бы три сборки, три сканирования и три места, где версии зависимостей могут разъехаться. Роль задаётся командой в compose.yaml.

ENTRYPOINT [] сброшен намеренно: команда задаётся целиком, и CMD служит значением по умолчанию для самой частой роли.

requirements.txt:

text
fastapi[standard]==0.141.1
pydantic-settings==2.14.2
psycopg[binary]==3.3.4
psycopg-pool==3.3.1
redis==8.1.0

requirements-dev.txt:

text
-r requirements.txt
pytest==9.1.1
pytest-cov==7.1.0
pytest-asyncio==1.4.0
httpx==0.28.1

pytest.ini:

ini
[pytest]
testpaths = tests
addopts = -q
asyncio_mode = auto
asyncio_default_fixture_loop_scope = function

.dockerignore:

text
.git
.gitignore
.venv
.v
__pycache__
*.pyc
.pytest_cache
.mypy_cache
.ruff_cache
htmlcov
.coverage
secrets/
compose*.yaml
README.md

Каталог secrets/ исключён из контекста сборки: без этого файл пароля попадёт в образ, а docker history покажет слой с ним.


Чего решение не делает

Стек не запускался. Docker на машине курса отсутствует. Ни docker compose up, ни сборка образа не выполнялись. Проверены отдельные части — миграции, хранилище, обработчик — против настоящего PostgreSQL, но их взаимодействие через Compose не проверено: ни порядок старта, ни healthcheck'и, ни internal: true, ни ограничения ресурсов.

Это существенная оговорка. Утверждения о поведении Compose в CHECKLIST.md — реконструкция по документации и по механизмам, разобранным в разделе 09.

Redis проверен. RedisQueue прогнан против настоящего Redis 8 в container'е: 6 тестов, включая семантику BLMOVE.

Приложение api не приведено целиком. Показаны части, которые меняются по сравнению с проектом 2: хранилище, секреты, миграции, обработчик, настройки, lifespan, Compose. Обработчики запросов, пробы и журналирование переносятся без изменений — в этом и был смысл выделения протокола Storage. Что именно меняется в settings.py и main.py — шаг 4a; пропуск этого шага виден сразу: api не стартует.

Резервное копирование не реализовано. Дополнительное задание 5 не выполнено. Оно важнее, чем кажется: том с данными есть, а процедуры восстановления нет, и это ровно тот случай, о котором предупреждает урок 19.1.

Время старта не измерено. Дополнительное задание 6 требует запуска стека.


Сравнение с вашей реализацией

Пять вопросов:

  1. Что делает ваш набор тестов, если базы нет? Если «проходит» — вы не знаете, проверено ли что-нибудь.
  2. Что произойдёт, если запустить миграции дважды одновременно?
  3. Может ли api стартовать раньше, чем создана схема?
  4. Что вы выбрали для очереди и почему? Оба варианта обоснованы, важен письменный довод.
  5. Проверяли ли вы старт десять раз подряд? Гонка проявляется не всегда.

Навигация

← Список проверок
← Техническое задание
Следующий проект: Production pipeline →
Вернуться к проектам
Главное оглавление

Markdown на GitHub ↗