Проект 3. Эталонное решение
Открывать после собственной реализации и прохождения CHECKLIST.md.
Решение — не единственно верное. Места, где обоснован другой выбор, отмечены явно.
Что проверено на самом деле
| Что | Как | Результат |
|---|---|---|
| Миграции | применение к настоящему PostgreSQL 16 | 3 применены, повтор ничего не делает |
| Отказ при правке применённой миграции | изменение файла и повторный запуск | отказ с указанием обеих контрольных сумм |
| Хранилище | 13 тестов против настоящей базы | все проходят |
| Белый список полей | попытки подстановки в GROUP BY | отклонены |
| Обработчик | 10 тестов, включая две копии сразу | SKIP LOCKED работает |
| Три состояния тестов | с базой, без базы, без базы с требованием | различимы |
| Весь набор | pytest | 37 пройдено, покрытие 80 % |
| Redis | 6 тестов против настоящего Redis 8 | все проходят |
| Стек целиком | docker compose up -d | пять сервисов, четыре healthy, migrate — Exited (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_url | UnknownSettingError, api в перезапусках |
create_app не читал secret | PoolTimeout: 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:
"""Чтение секретов из файлов.
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 в окружение всех трёх ролей и на этом останавливалась. Стек не поднимался вовсе:
$ 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:
-- Схема событий.
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:
-- Индексы под группировку. Отдельной миграцией, потому что первая
-- уже применена в эксплуатации и переписывать её нельзя.
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:
-- Признак обработки события фоновым обработчиком.
ALTER TABLE events ADD COLUMN processed_at TIMESTAMPTZ;
-- Частичный индекс: строк с NULL обычно мало, и запрос обработчика
-- «дай необработанные» должен читать только их.
CREATE INDEX events_unprocessed ON events (id) WHERE processed_at IS NULL;
Почему индексы отдельной миграцией. Первая уже применена в эксплуатации, и переписывать её нельзя: у тех, кто применил, изменения не появятся. Это не формальность — именно поэтому раннер отвергает изменённый файл.
migrations/runner.py:
"""Применение миграций.
Три свойства, каждое из которых проверено тестом:
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:
"""Хранилище событий в 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:
"""Очередь задач.
Протокол отделён от реализации по той же причине, что и хранилище
в проекте 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:
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:
"""Фоновый обработчик событий.
Забирает необработанные события из базы пачками и помечает их.
Останавливается по 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:
"""Приспособления для тестов с настоящей базой.
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: она превращает пропуск в отказ. Три состояния проверены запуском:
с базой: 37 passed
без базы: 4 passed, 33 skipped
без базы, REQUIRE_DB=1: 33 errors — «интеграционные тесты не выполнялись»
Это то же различение, что в уроке 16.4: «не проверено» должно отличаться от «проверено и чисто».
Схема пересоздаётся, а не очищается. TRUNCATE был бы быстрее, но тесты миграций меняют саму схему, и остаток от них испортил бы следующий тест.
tests/test_migrations.py:
"""Миграции: порядок, повторный запуск, изменение применённого файла."""
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:
"""Хранилище в 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:
"""Обработчик: пачки, остановка, взаимодействие с базой."""
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
Фактический результат:
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: два новых поля
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 это неизвестные имена:
$ 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: пул вместо словаря
@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() # пул закрывается ПОСЛЕ дренажа
Импорты меняются соответственно:
from psycopg_pool import AsyncConnectionPool
from .secrets import export_pgpassword
from .storage_pg import PostgresStorage
# MemoryStorage больше не нужен в рабочем коде — но остаётся в тестах
И create_app получает ту же строку, что миграции и обработчик:
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:
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:
# Конфигурация для тестов: 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:
# Конфигурация для разработки: 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:
# 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:
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:
-r requirements.txt
pytest==9.1.1
pytest-cov==7.1.0
pytest-asyncio==1.4.0
httpx==0.28.1
pytest.ini:
[pytest]
testpaths = tests
addopts = -q
asyncio_mode = auto
asyncio_default_fixture_loop_scope = function
.dockerignore:
.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 требует запуска стека.
Сравнение с вашей реализацией
Пять вопросов:
- Что делает ваш набор тестов, если базы нет? Если «проходит» — вы не знаете, проверено ли что-нибудь.
- Что произойдёт, если запустить миграции дважды одновременно?
- Может ли
apiстартовать раньше, чем создана схема? - Что вы выбрали для очереди и почему? Оба варианта обоснованы, важен письменный довод.
- Проверяли ли вы старт десять раз подряд? Гонка проявляется не всегда.
Навигация
← Список проверок
← Техническое задание
Следующий проект: Production pipeline →
Вернуться к проектам
Главное оглавление