diff --git a/README.md b/README.md index c09a8eb..ed8d2ca 100644 --- a/README.md +++ b/README.md @@ -98,7 +98,7 @@ docker compose up -d --build # UI: http://localhost:8000, OpenAPI: /docs python3 -m venv venv && ./venv/bin/pip install -r requirements.txt set -a; . ./.env; set +a ./venv/bin/uvicorn app.main:app --reload -./venv/bin/python -m pytest # 28 тестов, фоновый опрос в тестах выключен +./venv/bin/python -m pytest # 29 тестов, фоновый опрос в тестах выключен ``` ## API v1 @@ -166,3 +166,4 @@ curl -s -H "Authorization: Bearer $API_TOKEN" http://localhost:8000/api/v1/devic - `019-performance-scaling` — кэш списка бакета (`BACKUPS_CACHE_TTL`, «Обновить список»); обработчики без обращений к event loop — обычные функции (пул потоков FastAPI), запись статуса и тяжёлые операции с БД в фоне — через `asyncio.to_thread`; SQLite — WAL и `busy_timeout`; файловая блокировка БД — один процесс на БД. - `020-custom-select-menus` — выпадающие списки (фильтры, «Группа») в стиле меню действий «⋯»: прогрессивное улучшение в JS, нативный `select` остаётся в разметке и работает без JavaScript. - `021-security-hardening` — отказ старта при небезопасных секретах (`API_TOKEN`/`SESSION_SECRET`/`ADMIN_PASSWORD`/`SECRET_KEY`), блокировка входа в UI по IP клиента, `SESSION_COOKIE_SECURE`, безопасный `next` в редиректах (`/ui/move`, `/backups/delete-many`). +- `022-async-db-remainder` — остаток п. 9 ревью: оставшиеся синхронные обращения к БД в async-обработчиках (API, UI, `ops`, `backups.search`) — через `asyncio.to_thread`; `jobs.start_jobs` разделён на синхронную `_create_jobs` (в потоке) и `async start_jobs` (создаёт задачи и планирует их в event loop); регрессионный тест-линтер (AST-обход `app/`) не даёт синхронным обращениям к БД вернуться в async-код. diff --git a/app/api/v1.py b/app/api/v1.py index d87d907..54a1ab1 100644 --- a/app/api/v1.py +++ b/app/api/v1.py @@ -1,3 +1,4 @@ +import asyncio import json from datetime import date, datetime from typing import Literal @@ -167,51 +168,53 @@ def delete_device(device_id: str): @router.post("/devices/{device_id}/refresh", response_model=DeviceOut) async def refresh_device(device_id: str): - devices.get_device(device_id) + await asyncio.to_thread(devices.get_device, device_id) await ops.refresh_status(device_id) - return DeviceOut.of(devices.get_device(device_id)) + return DeviceOut.of(await asyncio.to_thread(devices.get_device, device_id)) @router.post("/devices/refresh", response_model=list[DeviceOut]) async def refresh_all(): await ops.refresh_many() - return [DeviceOut.of(d) for d in devices.list_devices()] + return [DeviceOut.of(d) for d in await asyncio.to_thread(devices.list_devices)] # --- операции (долгие — через jobs) --- @router.post("/devices/{device_id}/backups", status_code=202) async def create_backup(device_id: str): - return {"job_ids": jobs.start_jobs("backup", [device_id])} + return {"job_ids": await jobs.start_jobs("backup", [device_id])} @router.put("/devices/{device_id}/update/channel", response_model=DeviceOut) async def put_channel(device_id: str, body: ChannelIn): - devices.get_device(device_id) + await asyncio.to_thread(devices.get_device, device_id) await ops.set_channel(device_id, body.channel) - return DeviceOut.of(devices.get_device(device_id)) + return DeviceOut.of(await asyncio.to_thread(devices.get_device, device_id)) @router.post("/devices/{device_id}/update/install", status_code=202) async def install_update(device_id: str): - return {"job_ids": jobs.start_jobs("ros_update", [device_id])} + return {"job_ids": await jobs.start_jobs("ros_update", [device_id])} @router.post("/devices/{device_id}/firmware/upgrade", status_code=202) async def upgrade_firmware(device_id: str): - return {"job_ids": jobs.start_jobs("fw_update", [device_id])} + return {"job_ids": await jobs.start_jobs("fw_update", [device_id])} @router.post("/batch/{action}", status_code=202) async def batch(action: Literal["backup", "ros_update", "fw_update"], body: BatchIn): """Групповая операция над списком устройств.""" - return {"job_ids": jobs.start_jobs(action, body.resolve())} + device_ids = await asyncio.to_thread(body.resolve) + return {"job_ids": await jobs.start_jobs(action, device_ids)} @router.put("/batch/channel", status_code=202) async def batch_channel(body: BatchChannelIn): """Групповая смена канала — фоновыми задачами (как /batch/{backup|ros_update|fw_update}).""" - return {"job_ids": jobs.start_jobs("set_channel", body.resolve(), {"channel": body.channel})} + device_ids = await asyncio.to_thread(body.resolve) + return {"job_ids": await jobs.start_jobs("set_channel", device_ids, {"channel": body.channel})} # --- резервные копии в S3 --- @@ -226,7 +229,7 @@ async def list_backups( q: str = "", refresh: bool = False, # принудительно перечитать бакет, минуя кэш (BACKUPS_CACHE_TTL) ): - name = devices.get_device(device_id).name if device_id else "" + name = (await asyncio.to_thread(devices.get_device, device_id)).name if device_id else "" return await backups.list_backups(name, group, kind, date_from, date_to, q, refresh=refresh) diff --git a/app/services/backups.py b/app/services/backups.py index 4493557..bda34cb 100644 --- a/app/services/backups.py +++ b/app/services/backups.py @@ -118,8 +118,11 @@ async def search(device: str = "", group: str = "", kind: str = "", date_from: d group: "" — все, "none" — без группы (в т.ч. бэкапы удалённых устройств), иначе ID группы (grp_…). Даты — по времени изменения объекта (UTC), включительно. Каждому файлу сопоставлен backup_id. refresh — принудительно перечитать бакет, минуя кэш.""" - group_names = {g["id"]: g["name"] for g in groups.list_groups()} - device_group = {d.name: group_names.get(d.group_id) for d in devices.list_devices()} + def _group_maps(): + names = {g["id"]: g["name"] for g in groups.list_groups()} + return names, {d.name: names.get(d.group_id) for d in devices.list_devices()} + + group_names, device_group = await asyncio.to_thread(_group_maps) want_group = group_names.get(group) if ids.is_id(group, "grp") else None everything = await bucket_objects(refresh) diff --git a/app/services/jobs.py b/app/services/jobs.py index e7b1307..b208600 100644 --- a/app/services/jobs.py +++ b/app/services/jobs.py @@ -47,12 +47,10 @@ async def _run(job_id: str, job_type: str, device_id: str, params: dict) -> None await asyncio.to_thread(_finish, job_id, "failed", str(e)) -def start_jobs(job_type: str, device_ids: list[str], params: dict | None = None) -> list[str]: - """Создаёт задачи для списка устройств и запускает их в фоне. params передаются раннеру именованными - аргументами (device_id, **params) и попадают в data события job.created. Возвращает ID задач.""" +def _create_jobs(job_type: str, device_ids: list[str], params: dict) -> list[tuple[str, str]]: + """Синхронная часть start_jobs: строки задач и события job.created — в потоке.""" if job_type not in JOB_TYPES: raise ValueError(f"Неизвестный тип задачи: {job_type}") - params = params or {} for did in device_ids: ids.check(did, "dev") started = [] @@ -67,6 +65,14 @@ def start_jobs(job_type: str, device_ids: list[str], params: dict | None = None) events.record("job.created", "job", j.id, f"Задача {job_type} для {d.name} создана", device_id=did, job_id=j.id, data={"type": job_type, **params}, s=s) started.append((j.id, did)) + return started + + +async def start_jobs(job_type: str, device_ids: list[str], params: dict | None = None) -> list[str]: + """Создаёт задачи для списка устройств и запускает их в фоне. params передаются раннеру именованными + аргументами (device_id, **params) и попадают в data события job.created. Возвращает ID задач.""" + params = params or {} + started = await asyncio.to_thread(_create_jobs, job_type, device_ids, params) for jid, did in started: task = asyncio.create_task(_run(jid, job_type, did, params)) # наследует контекст (актор) _tasks.add(task) diff --git a/app/services/ops.py b/app/services/ops.py index 7361393..6b593cf 100644 --- a/app/services/ops.py +++ b/app/services/ops.py @@ -29,9 +29,13 @@ def _save_status(device_id: str, fields: dict) -> None: data={"error": (fields.get("last_error") or "")[:300]} if not up else None, s=s) +async def _conn(device_id: str) -> devices.Conn: + return await asyncio.to_thread(devices.get_conn, device_id) + + async def refresh_status(device_id: str) -> None: """Полный опрос: статус, версии и проверка обновлений. Недоступность — не исключение.""" - conn = await asyncio.to_thread(devices.get_conn, device_id) + conn = await _conn(device_id) try: async with devices.open_client(conn) as c: status = await ros.get_status(c) @@ -49,7 +53,7 @@ async def poll_device(device_id: str, full: bool = False) -> None: if full or d.online is not True or not d.status: return await refresh_status(device_id) try: - conn = await asyncio.to_thread(devices.get_conn, device_id) + conn = await _conn(device_id) async with devices.open_client(conn) as c: res = await c.get("system/resource") except RosError as e: @@ -59,7 +63,7 @@ async def poll_device(device_id: str, full: bool = False) -> None: async def refresh_many(device_ids: list[str] | None = None) -> None: - targets = device_ids if device_ids is not None else [d.id for d in devices.list_devices()] + targets = device_ids if device_ids is not None else [d.id for d in await asyncio.to_thread(devices.list_devices)] sem = asyncio.Semaphore(get_settings().max_concurrency) async def one(i: str): @@ -72,7 +76,7 @@ async def refresh_many(device_ids: list[str] | None = None) -> None: async def run_backup(device_id: str) -> str: """Бэкап .backup + .rsc: создать на устройстве → скачать через REST во временную папку → загрузить в S3 от имени сервера → удалить файлы с устройства.""" - conn = devices.get_conn(device_id) + conn = await _conn(device_id) stamp = datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S") base = f"ros_control-{stamp}" files = { # имя на устройстве -> ключ в S3 @@ -131,7 +135,7 @@ async def run_backup(device_id: str) -> str: async def set_channel(device_id: str, channel: str) -> None: - async with devices.open_client(devices.get_conn(device_id)) as c: + async with devices.open_client(await _conn(device_id)) as c: await ros.set_channel(c, channel) await refresh_status(device_id) @@ -143,10 +147,10 @@ async def run_set_channel(device_id: str, channel: str) -> str: async def run_ros_update(device_id: str) -> str: - async with devices.open_client(devices.get_conn(device_id)) as c: + async with devices.open_client(await _conn(device_id)) as c: return await ros.install_ros_update(c) async def run_fw_update(device_id: str) -> str: - async with devices.open_client(devices.get_conn(device_id)) as c: + async with devices.open_client(await _conn(device_id)) as c: return await ros.upgrade_firmware(c) diff --git a/app/ui/routes.py b/app/ui/routes.py index 36fa631..6320a5e 100644 --- a/app/ui/routes.py +++ b/app/ui/routes.py @@ -1,4 +1,5 @@ """WEB UI: Jinja2 + HTMX. Использует те же сервисы, что и JSON API.""" +import asyncio import json from datetime import date from pathlib import Path @@ -200,22 +201,22 @@ def jobs_table(request: Request): async def refresh(request: Request, device_ids: list[str] = Form(default=[])): flt = _flt(await request.form()) # без выбора обновляем все устройства текущего представления (группа + фильтры) - await ops.refresh_many(device_ids or [d.id for d in _visible(flt)]) - return _render(request, "_devices.html", oob=True, **_devices_ctx(flt)) + await ops.refresh_many(device_ids or [d.id for d in await asyncio.to_thread(_visible, flt)]) + return _render(request, "_devices.html", oob=True, **await asyncio.to_thread(_devices_ctx, flt)) @router.post("/ui/batch/{action}", response_class=HTMLResponse, dependencies=[Depends(require_login)]) async def batch(request: Request, action: str, device_ids: list[str] = Form(default=[]), channel: str = Form(default="")): if not device_ids: - return _render(request, "_jobs.html", jobs=jobs.list_jobs(15), flash="Выберите устройства") + return _render(request, "_jobs.html", jobs=await asyncio.to_thread(jobs.list_jobs, 15), flash="Выберите устройства") if action == "channel": if channel not in CHANNELS: raise ValueError(f"Неизвестный канал: {channel}") - jobs.start_jobs("set_channel", device_ids, {"channel": channel}) + await jobs.start_jobs("set_channel", device_ids, {"channel": channel}) else: - jobs.start_jobs(action, device_ids) - return _render(request, "_jobs.html", jobs=jobs.list_jobs(15)) + await jobs.start_jobs(action, device_ids) + return _render(request, "_jobs.html", jobs=await asyncio.to_thread(jobs.list_jobs, 15)) @router.post("/ui/devices/{device_id}/{action}", response_class=HTMLResponse, @@ -227,9 +228,10 @@ async def device_action(request: Request, device_id: str, action: str, channel: elif action == "channel": await ops.set_channel(device_id, channel) else: - jobs.start_jobs(action, [device_id]) - return _render(request, "_jobs.html", jobs=jobs.list_jobs(15)) - return _render(request, "_devices.html", oob=True, **_devices_ctx(_flt(await request.form()))) + await jobs.start_jobs(action, [device_id]) + return _render(request, "_jobs.html", jobs=await asyncio.to_thread(jobs.list_jobs, 15)) + flt = _flt(await request.form()) + return _render(request, "_devices.html", oob=True, **await asyncio.to_thread(_devices_ctx, flt)) def _safe_next(value: str, default: str = "/", prefix: str = "/") -> str: @@ -360,11 +362,15 @@ async def device_create(request: Request, name: str = Form(), host: str = Form() group_id: str = Form(""), new_group: str = Form(""), note: str = Form("")): v = dict(name=name, host=host, port=port, username=username, group_id=group_id, new_group=new_group, use_tls=use_tls, verify_tls=verify_tls, note=note) - try: + + def _create(): gid, ng = _group_choice(group_id, new_group) - d = devices.create_device(name, host, port, username, password, verify_tls, use_tls, gid, ng, note) + return devices.create_device(name, host, port, username, password, verify_tls, use_tls, gid, ng, note) + + try: + d = await asyncio.to_thread(_create) except ValueError as e: - return _device_form(request, None, v, str(e)) + return await asyncio.to_thread(_device_form, request, None, v, str(e)) await ops.refresh_status(d.id) return _done(request) @@ -420,12 +426,17 @@ async def backups_page(request: Request, device: str = "", group: str = "", kind error = str(e) flt = dict(device=device, group=group, kind=kind, date_from=date_from, date_to=date_to, q=q) non_empty = {k: v for k, v in flt.items() if v} + + def _lists(): + return devices.list_devices(), groups.list_groups() + + dev_list, grp_list = await asyncio.to_thread(_lists) return _render(request, "backups.html", section="backups", items=items, stats=stats, error=error, flt=flt, flt_active=sum(1 for v in flt.values() if v), bucket=get_settings().s3_bucket, page_url="/backups" + ("?" + urlencode(non_empty) if non_empty else ""), refresh_url="/backups?" + urlencode({**non_empty, "refresh": "1"}), deleted=deleted, failed=failed, - devices=devices.list_devices(), groups=groups.list_groups()) + devices=dev_list, groups=grp_list) @router.post("/backups/delete-many", dependencies=[Depends(require_login)]) diff --git a/docs/changes/022-async-db-remainder/plan.md b/docs/changes/022-async-db-remainder/plan.md new file mode 100644 index 0000000..fbee7f8 --- /dev/null +++ b/docs/changes/022-async-db-remainder/plan.md @@ -0,0 +1,81 @@ +# План: 022 — остаток п. 9: синхронная БД в async-коде + +## Context + +Ревью `docs/reviews/2026-09-28-1243-codebase-review.md`, п. 9 (частично закрыт в 019). В 019 обработчики без `await` стали `def`, +фоновые записи ушли в `asyncio.to_thread`, но в **async-функциях** остались прямые синхронные вызовы сервисов, которые ходят в SQLite +(`session_scope`) и блокируют event loop. В ревью названы два места; разведка (AST: async-функции, вызывающие функции с `session_scope` +напрямую или через одну ступень) нашла их больше: + +| Где | Синхронные вызовы БД | +|---|---| +| `app/services/backups.py::search` | `groups.list_groups`, `devices.list_devices` | +| `app/services/ops.py::run_backup`, `set_channel`, `run_ros_update`, `run_fw_update` | `devices.get_conn` (чтение + расшифровка пароля) | +| `app/services/ops.py::refresh_many` | `devices.list_devices` | +| `app/api/v1.py::refresh_device`, `put_channel`, `list_backups` | `devices.get_device` | +| `app/api/v1.py::refresh_all` | `devices.list_devices` | +| `app/api/v1.py::create_backup`, `install_update`, `upgrade_firmware`, `batch`, `batch_channel` | `jobs.start_jobs`, `BatchIn.resolve` | +| `app/ui/routes.py::refresh`, `device_action` | `_visible`, `_devices_ctx`, `jobs.list_jobs`, `jobs.start_jobs` | +| `app/ui/routes.py::batch` | `jobs.start_jobs`, `jobs.list_jobs` | +| `app/ui/routes.py::device_create` | `devices.create_device`, `_device_form` | +| `app/ui/routes.py::backups_page` | `devices.list_devices`, `groups.list_groups` | + +Ложные срабатывания: `events.record` внутри `ops.run_backup` — во вложенных `_create_row/_finish_row`, уже через `to_thread`. +`app/main.py::lifespan` (`init_db`, `fail_stale_jobs`) — выполняется до приёма запросов; **не меняем**. + +Решений пользователя не требуется: объём — «оставшиеся задачи п. 9», подход тот же, что в 019. + +## Изменения + +### Общий подход +- Синхронный вызов в async-функции → `await asyncio.to_thread(fn, *args)` (контекст копируется — актор и `job_id` в ContextVar сохраняются, как в 019). + Сигнатуры синхронных сервисов не меняются. +- Где несколько подряд идущих синхронных вызовов готовят данные для одного ответа (`_devices_ctx` + `list_jobs`, `list_devices` + `list_groups`) — + одна вложенная синхронная функция и один `to_thread`, а не серия переключений. + +### `app/services/jobs.py` — `start_jobs` +- Сейчас синхронная: создаёт строки задач в БД и затем `asyncio.create_task` (нужен работающий event loop — поэтому вызывающие обработчики остались `async`). +- Разделить: `_create_jobs(job_type, device_ids, params) -> list[tuple[job_id, device_id]]` — синхронная, только БД и события (с проверками + типа задачи и ID, как сейчас); `async def start_jobs(job_type, device_ids, params=None) -> list[str]` — `await asyncio.to_thread(_create_jobs, …)`, + затем `create_task` в event loop. Все вызовы `start_jobs` — с `await` (API, UI, тесты). + +### `app/services/ops.py` +- Хелпер `async def _conn(device_id) -> devices.Conn` = `await asyncio.to_thread(devices.get_conn, device_id)`; использовать в + `refresh_status`, `poll_device` (уже через `to_thread` — заменить на хелпер для единообразия), `run_backup`, `set_channel`, `run_ros_update`, `run_fw_update`. +- `refresh_many`: `devices.list_devices` → `to_thread`. + +### `app/services/backups.py::search` +- `groups.list_groups()` и `devices.list_devices()` — одним `to_thread` (вложенная функция, возвращающая `group_names`, `device_group`). + +### `app/api/v1.py` +- `refresh_device`, `put_channel`, `list_backups`: `devices.get_device` → `to_thread`. +- `refresh_all`: итоговый `devices.list_devices` → `to_thread`. +- `batch`, `batch_channel`: `body.resolve()` → `to_thread`; `await jobs.start_jobs(...)`. `create_backup`, `install_update`, `upgrade_firmware` — `await jobs.start_jobs(...)`. + +### `app/ui/routes.py` +- `refresh`, `batch`, `device_action`: подготовка контекста (`_visible`, `_devices_ctx`, `jobs.list_jobs`) — через `to_thread`; `await jobs.start_jobs(...)`. + `await request.form()` остаётся в event loop (до `to_thread`). +- `device_create`: `devices.create_device` и повторный рендер формы с ошибкой (`_device_form` читает группы) — через `to_thread`; + `await ops.refresh_status` без изменений. +- `backups_page`: `devices.list_devices`, `groups.list_groups` — через `to_thread`. + +## Тесты (минимально) +- **Регрессионный тест-линтер** (`tests/test_app.py`): AST-обход `app/` — в `async def` нет прямых вызовов функций, использующих + `session_scope` (напрямую или через одну ступень вызовов), кроме разрешённого списка (`lifespan`); вызовы, переданные ссылкой в + `asyncio.to_thread(fn, …)`, и вложенные `def` не считаются. Логика — как у скрипта разведки из Context. Это закрепляет п. 9 от регресса. +- Существующие тесты, вызывающие `jobs.start_jobs`, — перевести на `await`. Актор в событиях задач (`api`, `ui:`) должен сохраниться — + подтверждается существующими тестами журнала. + +## Документация +README: строка 022 в «История изменений», число тестов. `summary.md` — оркестратор. + +## Исполнение +Исполнитель (Sonnet): код, тест, README, пересборка стенда. Тесты не запускает, не коммитит, `.env` не читает. + +## Проверка +- `pytest` — все зелёные; тест-линтер падает, если временно вернуть прямой вызов в async-функцию (мутационная проверка — оркестратор). +- Повторный запуск скрипта разведки — пусто (кроме `lifespan`). +- Стенд (override 8001, `--force-recreate`): `/login` 200, новый код в контейнере; сценарий: временное устройство `192.0.2.1` → + `POST /api/v1/devices/{id}/backups` → 202, задача `backup` завершается `failed` (недоступно) с актором `api` в событиях; `PUT /batch/channel` → 202; + параллельный `GET /api/v1/devices` во время `refresh` отвечает сразу; устройство удаляется. Боевые данные — сверка по ID. +- Ручная проверка UI — пользователь: «Обновить статус», групповые действия, меню «⋯» устройства, добавление устройства, страница «Бэкапы». diff --git a/docs/changes/022-async-db-remainder/summary.md b/docs/changes/022-async-db-remainder/summary.md new file mode 100644 index 0000000..8d8f11c --- /dev/null +++ b/docs/changes/022-async-db-remainder/summary.md @@ -0,0 +1,25 @@ +# Итоги: 022 — остаток п. 9: синхронная БД в async-коде + +Источник — ревью `docs/reviews/2026-09-28-1243-codebase-review.md`, п. 9 (частично закрыт в 019). + +## Сделано +- Разведка по AST нашла 14 async-функций с прямыми синхронными обращениями к SQLite (в ревью были названы 2) — все переведены на `asyncio.to_thread`; + исключение — `lifespan` (выполняется до приёма запросов). +- `jobs.start_jobs` — `async`: БД и события `job.created` — синхронная `_create_jobs` в потоке, `create_task` — в event loop; все вызовы с `await`. +- `ops._conn` — чтение устройства и расшифровка пароля в потоке во всех операциях (опрос, бэкап, канал, обновления ROS/FW); `refresh_many` — список устройств в потоке. +- Группировка нескольких запросов одного ответа в один `to_thread`: `backups.search`, `backups_page`, `device_create`. +- API (`refresh_device`, `put_channel`, `refresh_all`, `list_backups`, `batch`, `batch_channel`, запуск задач) и UI (`refresh`, `batch`, `device_action`, + `device_create`, `backups_page`) — без синхронной БД в event loop. +- Тест-линтер `test_no_sync_db_calls_in_async_functions`: AST-обход `app/` — в `async def` нет прямых вызовов функций с `session_scope` (напрямую или через одну ступень); + ссылки в `asyncio.to_thread(fn, …)` и вложенные `def` не считаются; сообщение называет файл, функцию и вызов. + +## Проверено +- `pytest`: 29 из 29. +- Мутация: прямой `get_device` в `refresh_device` → линтер падает с точным местом; файл восстановлен. +- Стенд (8001): бэкап и групповая смена канала для временного устройства → 202, задачи `failed` (`ConnectTimeout`, ожидаемо), актор всех событий задач — `api` + (контекст в потоке сохраняется); во время 4-секундного `refresh` параллельные `GET /devices` — 0,011 с, `GET /backups` — 0,125 с. +- Боевые данные: группы 3/3, устройства 14/14; добавлены только данные проверки (2 задачи, запись неудачного бэкапа, события); временное устройство удалено. + +## Оговорки +- Линтер видит одну ступень транзитивности: функция, обращающаяся к БД через две и более промежуточных, не будет обнаружена. +- Ручная проверка UI пользователем на момент коммита не подтверждена. diff --git a/docs/reviews/2026-09-28-1243-codebase-review.md b/docs/reviews/2026-09-28-1243-codebase-review.md new file mode 100644 index 0000000..75dcd58 --- /dev/null +++ b/docs/reviews/2026-09-28-1243-codebase-review.md @@ -0,0 +1,79 @@ +# Ревью кодовой базы ros_control — 2026-09-28 12:43 MSK + +Повторное ревью. Состояние на коммит `c407972`. Тесты: 28/28 проходят. +Предыдущее ревью: [`2026-09-27-codebase-review.md`](2026-09-27-codebase-review.md) (коммит `26dd1f8`). + +## Итог + +Из 12 замечаний предыдущего ревью закрыто 10 (одно — частично), открыто 2. Закрыты разделы «Безопасность», +«Корректность и согласованность», «Производительность и масштабирование»; раздел «Эксплуатация» не менялся. + +| Изменение | Коммит | Пункты | +|---|---|---| +| `018-correctness-consistency` | `123b5ab` | 5, 6, 7 | +| `019-performance-scaling` | `ba8ac7b` | 8, 9 (частично), 10 | +| `020-custom-select-menus` | `6b5c521` | — (доработка UI по запросу) | +| `021-security-hardening` | `c407972` | 1, 2, 3, 4 | + +## Статус замечаний предыдущего ревью + +### Безопасность + +| № | Замечание | Приоритет | Статус | Подтверждение в коде | +|---|---|---|---|---| +| 1 | Небезопасные значения секретов по умолчанию | Высокий | ✅ 021 | `app/config.py::insecure_settings`; проверка первым шагом в `app/main.py::lifespan` — при пустых, `change-me` или коротких секретах (`API_TOKEN`, `SESSION_SECRET` ≥ 32, `ADMIN_PASSWORD` ≥ 12) и невалидном `SECRET_KEY` приложение не стартует | +| 2 | Вход в UI без защиты от перебора | Высокий | ✅ 021 | `app/ui/routes.py::login`: 5 неверных за 10 минут → IP заблокирован на 10 минут (429, пароль не проверяется); `auth.locked` пишется один раз | +| 3 | Сессионная cookie без флага Secure | Средний | ✅ 021 | `SESSION_COOKIE_SECURE` → `SessionMiddleware(https_only=…)`; по умолчанию выключен, за TLS — включить | +| 4 | Открытый редирект в `/ui/move` | Низкий | ✅ 021 | `app/ui/routes.py::_safe_next` — только локальный путь; применён в `/ui/move` и групповом удалении бэкапов | + +### Корректность и согласованность + +| № | Замечание | Приоритет | Статус | Подтверждение в коде | +|---|---|---|---|---| +| 5 | Одиночное удаление бэкапа в обход сервиса | Средний | ✅ 018 | `app/ui/routes.py::backup_delete` → `backups.delete_many([key])` | +| 6 | Две системы миграций | Низкий | ✅ 018 | `db._migrate` удалён; колонки старой схемы — `migrations._add_legacy_columns` при `user_version < 1` | +| 7 | Синхронная групповая смена канала | Средний | ✅ 018 | Задачи `set_channel`; `PUT /api/v1/batch/channel` → 202 `{"job_ids"}` | + +### Производительность и масштабирование + +| № | Замечание | Приоритет | Статус | Подтверждение в коде | +|---|---|---|---|---| +| 8 | Каждое открытие «Бэкапов» перечитывает весь бакет | Средний | ✅ 019 | `app/services/backups.py::bucket_objects` — кэш на `BACKUPS_CACHE_TTL`, счётчик поколений, `invalidate()` после бэкапа и удаления, `refresh=1` | +| 9 | Синхронная БД в async-обработчиках | Низкий | 🟡 019, частично | 44 обработчика стали `def`; фоновые записи — `asyncio.to_thread`; SQLite WAL + `busy_timeout`. Остаток — см. «Открытые замечания», п. 9 | +| 10 | Состояние в памяти процесса | Низкий | ✅ 019 | `app/process_lock.py` — `flock` на `<файл БД>.lock` в `lifespan`; README «Ограничения» | + +### Эксплуатация + +| № | Замечание | Приоритет | Статус | Подтверждение в коде | +|---|---|---|---|---| +| 11 | Нет healthcheck, логирование не настроено | Низкий | ❌ Открыт | Нет `HEALTHCHECK` в `Dockerfile` и `docker-compose.yml`; логирование нигде не настраивается | +| 12 | Тесты в одном файле | Низкий | ❌ Открыт, растёт | `tests/test_app.py`: 570 → 753 строки, 20 → 28 тестов | + +## Открытые замечания + +| № | Замечание | Где | Приоритет | +|---|---|---|---| +| 9 | Остаток синхронной работы с БД в async-коде: `groups.list_groups()` и `devices.list_devices()` в async-обработчике страницы «Бэкапы»; `devices.get_conn()` (чтение и расшифровка пароля) | `app/services/backups.py::search`, `app/services/ops.py::run_backup` | Низкий | +| 11 | Нет healthcheck; логирование не настроено | `Dockerfile`, `docker-compose.yml`, `app/main.py` | Низкий | +| 12 | Тесты в одном файле, файл растёт | `tests/test_app.py` | Низкий | + +## Новые замечания + +| № | Замечание | Где | Приоритет | +|---|---|---|---| +| 13 | Копирование БД в режиме WAL не описано: копия файла БД без `-wal` может быть неполной; нужен `sqlite3 .backup` или копирование вместе с `-wal` | `README.md` | Средний | +| 14 | За reverse-proxy блокировка входа видит один IP прокси (нет настройки доверенных прокси и разбора `X-Forwarded-For`); ограничение описано в README | `app/ui/routes.py::login` | Низкий (прокси сейчас нет) | +| 15 | `refresh=1` остаётся в адресе страницы «Бэкапы» после «Обновить список» — повторные перезагрузки идут мимо кэша | `app/ui/routes.py::backups_page` | Низкий | +| 16 | Стенд запускается на порту 8001 через override-файл вне репозитория: порт 8000 на хосте занят посторонним процессом (`telemetry_web`); штатный `docker-compose.yml` без override не стартует. Порт не параметризован | `docker-compose.yml`, окружение хоста | Низкий (окружение) | + +## Проверка UI + +Ручная проверка интерфейса пользователем по изменениям 019–021 не подтверждена (отмечено в `summary.md` и сообщениях коммитов). +Особенно важно для 020 (выпадающие списки): автотестов JS в проекте нет. + +## Рекомендуемый порядок + +1. **Изменение 022:** п. 11 (`HEALTHCHECK`, базовая настройка логов), остаток п. 9, п. 13 (копирование БД в README), п. 15. +2. П. 16 — параметризовать порт в `docker-compose.yml` (`${HTTP_PORT:-8000}`) или освободить порт 8000 на хосте. +3. П. 12 — разбить тесты по модулям отдельным изменением (только структура тестов). +4. П. 14 — при появлении reverse-proxy: доверенные прокси и реальный IP клиента. diff --git a/tests/test_app.py b/tests/test_app.py index 33b11ad..e9e5534 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -1,9 +1,11 @@ +import ast import asyncio import html import json import re import sqlite3 from datetime import date, datetime, timedelta, timezone +from pathlib import Path import httpx import pytest @@ -478,7 +480,7 @@ async def test_events_link_entities(monkeypatch): return "готово" monkeypatch.setitem(jobs.JOB_TYPES, "backup", fake_backup) - [job_id] = jobs.start_jobs("backup", [d.id]) + [job_id] = await jobs.start_jobs("backup", [d.id]) await asyncio.sleep(0.3) evs = events.list_events(device_id=d.id) @@ -751,3 +753,69 @@ def test_move_redirect_rejects_open_redirect_next(): assert r.headers["location"] == "/" r = c.post("/ui/move", data={"next": "/?f_group=none"}, follow_redirects=False) assert r.headers["location"] == "/?f_group=none" # обычный путь с фильтром сохраняется + + +def test_no_sync_db_calls_in_async_functions(): + """Регресс п.9 ревью: в async-функциях app/ не должно быть прямых вызовов синхронных функций, + обращающихся к БД (session_scope) — напрямую или через один уровень вызовов, — это блокирует event loop. + Правильный способ — asyncio.to_thread(fn, ...). Ссылка на функцию, переданная в asyncio.to_thread(fn, ...), + вызовом не считается; тела вложенных def/async def в область видимости внешней функции не входят. + Единственное исключение — app.main.lifespan (выполняется до приёма запросов).""" + app_dir = Path(__file__).resolve().parent.parent / "app" + allowed = {"lifespan"} + + class _OwnScope(ast.NodeVisitor): + """Вызовы прямо в теле функции, не заходя в тела вложенных def/async def.""" + + def __init__(self): + self.calls: list[ast.Call] = [] + + def visit_FunctionDef(self, node): # не спускаемся во вложенную функцию + pass + + def visit_AsyncFunctionDef(self, node): + pass + + def visit_Call(self, node): + self.calls.append(node) + self.generic_visit(node) + + def own_calls(func) -> list[ast.Call]: + v = _OwnScope() + for stmt in func.body: + v.visit(stmt) + return v.calls + + def call_name(call: ast.Call) -> str | None: + f = call.func + if isinstance(f, ast.Name): + return f.id + if isinstance(f, ast.Attribute): + return f.attr + return None + + sync_funcs, async_funcs = [], [] + for path in sorted(app_dir.rglob("*.py")): + tree = ast.parse(path.read_text(), filename=str(path)) + for node in ast.walk(tree): + if isinstance(node, ast.AsyncFunctionDef): + async_funcs.append((path, node)) + elif isinstance(node, ast.FunctionDef): + sync_funcs.append((path, node)) + + # функции, обращающиеся к БД напрямую, и функции, вызывающие их (один уровень) + level0 = {f.name for _, f in sync_funcs if any(call_name(c) == "session_scope" for c in own_calls(f))} + level1 = {f.name for _, f in sync_funcs + if f.name not in level0 and any(call_name(c) in level0 for c in own_calls(f))} + dangerous = level0 | level1 + + violations = [] + for path, func in async_funcs: + if func.name in allowed: + continue + for c in own_calls(func): + name = call_name(c) + if name in dangerous: + violations.append(f"{path.relative_to(app_dir.parent)}::{func.name} вызывает {name}() напрямую, " + "в обход asyncio.to_thread — обращение к БД блокирует event loop") + assert not violations, "\n".join(violations)