diff --git a/backend/app/api/auth.py b/backend/app/api/auth.py index f6e1941..0570d32 100644 --- a/backend/app/api/auth.py +++ b/backend/app/api/auth.py @@ -28,7 +28,7 @@ from argon2 import PasswordHasher from argon2.exceptions import VerifyMismatchError from fastapi import APIRouter, HTTPException, Request, WebSocket from pydantic import BaseModel, Field -from sqlalchemy import select +from sqlalchemy import select, update from sqlalchemy.exc import IntegrityError from starlette.websockets import WebSocketDisconnect @@ -57,8 +57,21 @@ AUTH_GENERATION_SYNC_SECONDS = 1.0 AUTH_GENERATION_MAX_AGE_SECONDS = 2.0 _generations_synced_at: float | None = None # Один lock на ещё неизвестный узлу логин: параллельные HTTP и WS входы -# разделяют один запрос к БД до очередной синхронизации (lct-42). -_lookup_locks: dict[str, asyncio.Lock] = {} +# разделяют один запрос к БД до очередной синхронизации (lct-42). Запись +# удаляется, когда lock больше никто не ждёт: иначе словарь растёт с каждым +# логином, который узел когда-либо проверял. +_lookup_locks: dict[str, "_LoginLookup"] = {} +# Логины, которых разовая проверка не нашла в users. Держатся до следующей +# сверки: повтор старой cookie не должен давать SELECT на каждый запрос. +_missing_logins: set[str] = set() + + +class _LoginLookup: + __slots__ = ("lock", "users") + + def __init__(self) -> None: + self.lock = asyncio.Lock() + self.users = 0 class _AuthStateUnavailable(Exception): @@ -105,6 +118,7 @@ def prime_generations(values: dict[str, int]) -> None: global _generations_synced_at _generations.clear() _generations.update(values) + _missing_logins.clear() _generations_synced_at = time.monotonic() @@ -119,7 +133,27 @@ async def sync_generations() -> None: known_before_query = set(_generations) async with get_sessionmaker()() as db: rows = (await db.execute(select(User.login, User.auth_version))).all() - current = {login: version for login, version in rows} + current = {login: version for login, version in rows} + # Версия узла выше БД, если отзыв не удалось записать (неудачный + # logout) или учётку пересоздали после удаления. Без записи в БД + # узлы расходятся навсегда: соседний узел выдаёт cookie со старой + # версией, а этот её отвергает. Поднимаем БД до версии узла — отзыв + # сохраняется, остальные узлы догоняют за одну сверку. При гонке со + # снимком БД уже выше, и UPDATE ничего не меняет. + ahead = { + login: _generations[login] + for login, version in current.items() + if login in _generations and _generations[login] > version + } + if ahead: + try: + for login, local in ahead.items(): + await db.execute(_raise_auth_version(login, local)) + await db.commit() + except Exception as exc: # noqa: BLE001 — повторим на следующей сверке + log.error("не удалось записать версию полномочий узла (%s)", + type(exc).__name__) + await db.rollback() for login, version in current.items(): previous = _generations.get(login) # Снимок БД мог устареть за время запроса; меньшая версия не должна @@ -142,40 +176,81 @@ async def sync_generations() -> None: for login in known_before_query - current.keys() - synthetic - _vanished: invalidate_login(login) _vanished.add(login) + # Снимок только что прочитан: учётка, созданная до него, уже в кэше. + _missing_logins.clear() global _generations_synced_at _generations_synced_at = time.monotonic() +def _raise_auth_version(login: str, version: int): + # Только вверх: параллельная запись в БД могла уже поднять версию выше. + return ( + update(User) + .where(User.login == login, User.auth_version < version) + .values(auth_version=version) + ) + + async def _resolve_unknown_login(login: str) -> int | None: """Проверить неизвестный узлу логин сразу, не дожидаясь опроса БД. Отозванный логин остаётся в `_generations` (см. `_vanished`), поэтому отсутствие в кэше означает, что этот узел ещё не видел учётку. Возвращает версию или None, если строки в `users` нет. При недоступной - БД вызывает `_AuthStateUnavailable`, чтобы вход был закрыт с 503/1013. + или зависшей БД вызывает `_AuthStateUnavailable`, чтобы вход был закрыт + с 503/1013. """ - lock = _lookup_locks.setdefault(login, asyncio.Lock()) - async with lock: - cached = _generations.get(login) - if cached is not None: - return cached # параллельный запрос уже получил версию - try: - async with get_sessionmaker()() as db: - version = await db.scalar( - select(User.auth_version).where(User.login == login) + if login in _missing_logins: + return None + entry = _lookup_locks.get(login) + if entry is None: + entry = _lookup_locks[login] = _LoginLookup() + entry.users += 1 + try: + # Тот же предел, что у сверки: при partition handshake не ждёт + # таймаута TCP, а очередь за lock не растягивает ожидание сверх него. + async with asyncio.timeout(AUTH_GENERATION_MAX_AGE_SECONDS): + async with entry.lock: + cached = _generations.get(login) + if cached is not None: + return cached # параллельный запрос уже получил версию + if login in _missing_logins: + return None + async with get_sessionmaker()() as db: + version = await db.scalar( + select(User.auth_version).where(User.login == login) + ) + # Пока шёл SELECT, локальный отзыв или синхронизация могли + # записать новую версию. Старый ответ не должен вернуть + # отозванную cookie. + cached = _generations.get(login) + if cached is not None: + return cached + if version is None: + _missing_logins.add(login) + return None + _generations[login] = version + log.warning( + "учётка %s найдена разовой проверкой до сверки узла", + login_log_marker(login), ) - except Exception as exc: # noqa: BLE001 — ошибка БД должна закрыть вход - log.error("разовая проверка полномочий не удалась (%s)", type(exc).__name__) - raise _AuthStateUnavailable from exc - # Пока шёл SELECT, локальный отзыв или синхронизация могли записать - # новую версию. Старый ответ не должен вернуть отозванную cookie. - cached = _generations.get(login) - if cached is not None: - return cached - if version is None: - return None - _generations[login] = version - return version + return version + except Exception as exc: # noqa: BLE001 — ошибка или таймаут БД закрывают вход + log.error("разовая проверка полномочий не удалась (%s)", type(exc).__name__) + raise _AuthStateUnavailable from exc + finally: + entry.users -= 1 + if entry.users == 0 and _lookup_locks.get(login) is entry: + del _lookup_locks[login] + + +def login_log_marker(login: str) -> str: + """Метка логина для журнала без самого логина. + + По ней кластерный смоук доказывает, что узел прошёл через разовую + проверку, а не увидел учётку обычной сверкой. + """ + return hashlib.sha256(f"lct-login:{login}".encode()).hexdigest()[:12] async def watch_generations() -> None: @@ -590,7 +665,20 @@ async def login(payload: LoginIn, request: Request) -> dict: await audit_required(user.login, user.role, "login.blocked") raise HTTPException(status_code=403, detail="blocked") - if _generations.get(user.login) != user.auth_version: + local_version = _generations.get(user.login) + if local_version is not None and local_version > user.auth_version: + # Отзыв на этом узле не дошёл до БД (неудачный logout). Сброс к версии + # БД вернул бы силу cookie, выданной до выхода; выдача cookie с + # версией узла без записи в БД развела бы узлы. Поднимаем БД. + try: + async with get_sessionmaker()() as db: + await db.execute(_raise_auth_version(user.login, local_version)) + await db.commit() + except Exception as exc: # noqa: BLE001 — cookie без записанной версии не выдаём + log.error("вход: версию полномочий не удалось записать (%s)", + type(exc).__name__) + raise HTTPException(status_code=503, detail="auth_state_unavailable") from exc + elif local_version != user.auth_version: invalidate_login(user.login, user.auth_version) who = Principal(