From b6617d7df8ab6fa06942034701fff51c21b59cae Mon Sep 17 00:00:00 2001 From: kaifarikman Date: Sun, 27 Sep 2026 17:05:29 +0300 Subject: [PATCH 1/7] =?UTF-8?q?fix:=20=D0=BD=D0=B5=D0=B8=D0=B7=D0=B2=D0=B5?= =?UTF-8?q?=D1=81=D1=82=D0=BD=D1=8B=D0=B9=20=D1=83=D0=B7=D0=BB=D1=83=20?= =?UTF-8?q?=D0=BB=D0=BE=D0=B3=D0=B8=D0=BD=20=D0=BF=D1=80=D0=BE=D0=B2=D0=B5?= =?UTF-8?q?=D1=80=D1=8F=D0=B5=D1=82=D1=81=D1=8F=20=D1=80=D0=B0=D0=B7=D0=BE?= =?UTF-8?q?=D0=B2=D0=BE,=20=D0=BD=D0=B5=20=D0=BE=D1=82=D0=B7=D1=8B=D0=B2?= =?UTF-8?q?=D0=B0=D0=B5=D1=82=D1=81=D1=8F=20=D1=81=D1=80=D0=B0=D0=B7=D1=83?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit lct-42: узел кластера принимал вход другого узла первую секунду за отзыв и рвал сокет FORBIDDEN. Неизвестный логин теперь проверяется разовым SELECT auth_version, результат кэшируется, параллельные запросы одного логина дедуплицируются; логина нет в users или БД недоступна — как раньше. --- backend/app/api/auth.py | 72 +++++++-- backend/tests/test_auth_hardening.py | 215 +++++++++++++++++++++++++++ 2 files changed, 275 insertions(+), 12 deletions(-) diff --git a/backend/app/api/auth.py b/backend/app/api/auth.py index 17af9ce..bb25c65 100644 --- a/backend/app/api/auth.py +++ b/backend/app/api/auth.py @@ -56,6 +56,14 @@ _vanished: set[str] = set() AUTH_GENERATION_SYNC_SECONDS = 1.0 AUTH_GENERATION_MAX_AGE_SECONDS = 2.0 _generations_synced_at: float | None = None +# One lock per login this node has not synced yet, so a concurrent HTTP +# request and WS handshake for the same just-migrated login collapse into a +# single SELECT instead of one each (lct-42). +_lookup_locks: dict[str, asyncio.Lock] = {} + + +class _AuthStateUnavailable(Exception): + """The one-shot lookup for a login unknown to this node could not reach PostgreSQL.""" async def _close_revoked(ws: WebSocket) -> None: @@ -134,6 +142,36 @@ async def sync_generations() -> None: _generations_synced_at = time.monotonic() +async def _resolve_unknown_login(login: str) -> int | None: + """Look up a login absent from this node's cache without waiting a tick. + + A revoked login stays cached in `_generations` (see `_vanished` above), + so absence from `_generations` unambiguously means "this node has not + synced this login yet" — an unrecognized login is checked directly + rather than treated as revoked. Returns the current epoch, or `None` if + the login is not (or no longer) in `users`. Raises + `_AuthStateUnavailable` if PostgreSQL cannot be reached; the caller + fails closed. + """ + lock = _lookup_locks.setdefault(login, asyncio.Lock()) + async with lock: + cached = _generations.get(login) + if cached is not None: + return cached # a concurrent request already resolved it + try: + async with get_sessionmaker()() as db: + version = await db.scalar( + select(User.auth_version).where(User.login == login) + ) + except Exception as exc: # noqa: BLE001 — fail closed, not FORBIDDEN + log.error("разовая проверка полномочий не удалась (%s)", type(exc).__name__) + raise _AuthStateUnavailable from exc + if version is None: + return None + _generations[login] = version + return version + + async def watch_generations() -> None: """Poll PostgreSQL once per node so remote logout/role changes close WS.""" while True: @@ -151,6 +189,21 @@ async def watch_generations() -> None: await asyncio.sleep(AUTH_GENERATION_SYNC_SECONDS) +async def _send_auth_state_unavailable(scope, send) -> None: + if scope["type"] == "websocket": + await send({"type": "websocket.close", "code": 1013}) + else: + await send({ + "type": "http.response.start", + "status": 503, + "headers": [(b"content-type", b"application/json")], + }) + await send({ + "type": "http.response.body", + "body": b'{"detail":"auth_state_unavailable"}', + }) + + class AuthVersionMiddleware: """Check signed-cookie epochs against the fresh, DB-synchronized node cache.""" @@ -186,22 +239,17 @@ class AuthVersionMiddleware: synced_at is None or time.monotonic() - synced_at > AUTH_GENERATION_MAX_AGE_SECONDS ): - if scope["type"] == "websocket": - await send({"type": "websocket.close", "code": 1013}) - else: - await send({ - "type": "http.response.start", - "status": 503, - "headers": [(b"content-type", b"application/json")], - }) - await send({ - "type": "http.response.body", - "body": b'{"detail":"auth_state_unavailable"}', - }) + await _send_auth_state_unavailable(scope, send) return cookie_version = session.get("auth_generation") version = _generations.get(login) + if version is None: + try: + version = await _resolve_unknown_login(login) + except _AuthStateUnavailable: + await _send_auth_state_unavailable(scope, send) + return if version is None or cookie_version != version: if session is not None: session.clear() diff --git a/backend/tests/test_auth_hardening.py b/backend/tests/test_auth_hardening.py index eec963a..00cd8bf 100644 --- a/backend/tests/test_auth_hardening.py +++ b/backend/tests/test_auth_hardening.py @@ -410,6 +410,221 @@ def test_auth_middleware_fails_closed_when_generation_cache_is_stale_but_allows_ asyncio.run(run()) +def test_middleware_resolves_and_caches_login_unknown_to_this_node(monkeypatch): + """Cluster handshake (lct-42): a login another node just authenticated is + fetched via a single SELECT rather than being treated as revoked, and the + result is cached so a second request for it does not query again.""" + login = "peer-node-fresh-login" + calls = [] + + class FakeDb: + async def scalar(self, _query): + calls.append(1) + return 5 + + class FakeSession: + async def __aenter__(self): + return FakeDb() + + async def __aexit__(self, *_args): + return None + + class InnerApp: + def __init__(self): + self.called = 0 + + async def __call__(self, _scope, _receive, _send): + self.called += 1 + + async def run(): + auth.prime_generations({}) + + def make_scope(): + return { + "type": "http", "path": "/api/admin/users", + "session": { + "principal": {"login": login}, + "auth_instance": auth._INSTANCE, + "auth_generation": 5, + }, + } + + async def receive(): + return {"type": "http.request", "body": b"", "more_body": False} + + async def send(_message): + return None + + first = InnerApp() + await auth.AuthVersionMiddleware(first)(make_scope(), receive, send) + assert first.called == 1, "a fresh, valid epoch must reach the route" + assert auth._generations[login] == 5 + assert calls == [1] + + second = InnerApp() + await auth.AuthVersionMiddleware(second)(make_scope(), receive, send) + assert second.called == 1 + assert calls == [1], "cached epoch must not trigger a second SELECT" + + settings = auth.get_settings().model_copy(update={"demo_no_db": False}) + monkeypatch.setattr(auth, "get_settings", lambda: settings) + monkeypatch.setattr(auth, "get_sessionmaker", lambda: lambda: FakeSession()) + try: + asyncio.run(run()) + finally: + auth._generations.pop(login, None) + auth._lookup_locks.pop(login, None) + + +def test_middleware_rejects_login_missing_from_users_via_one_shot_query(monkeypatch): + login = "peer-node-deleted-login" + + class FakeDb: + async def scalar(self, _query): + return None + + class FakeSession: + async def __aenter__(self): + return FakeDb() + + async def __aexit__(self, *_args): + return None + + class InnerApp: + def __init__(self): + self.called = False + + async def __call__(self, _scope, _receive, _send): + self.called = True + + async def run(): + auth.prime_generations({}) + cookie_session = { + "principal": {"login": login}, + "auth_instance": auth._INSTANCE, + "auth_generation": 3, + } + scope = {"type": "http", "path": "/api/admin/users", "session": cookie_session} + + async def receive(): + return {"type": "http.request", "body": b"", "more_body": False} + + async def send(_message): + return None + + protected = InnerApp() + await auth.AuthVersionMiddleware(protected)(scope, receive, send) + assert protected.called, "route still runs but the session was cleared below" + assert cookie_session == {} + assert login not in auth._generations + + settings = auth.get_settings().model_copy(update={"demo_no_db": False}) + monkeypatch.setattr(auth, "get_settings", lambda: settings) + monkeypatch.setattr(auth, "get_sessionmaker", lambda: lambda: FakeSession()) + try: + asyncio.run(run()) + finally: + auth._lookup_locks.pop(login, None) + + +def test_unknown_login_lookup_deduplicates_concurrent_requests_into_one_select(monkeypatch): + login = "peer-node-concurrent-login" + calls = [] + + class FakeDb: + async def scalar(self, _query): + calls.append(1) + await asyncio.sleep(0.01) # widen the window for a racing second caller + return 9 + + class FakeSession: + async def __aenter__(self): + return FakeDb() + + async def __aexit__(self, *_args): + return None + + async def run(): + auth.prime_generations({}) + results = await asyncio.gather( + auth._resolve_unknown_login(login), + auth._resolve_unknown_login(login), + auth._resolve_unknown_login(login), + ) + assert results == [9, 9, 9] + assert calls == [1], "concurrent lookups for one login must issue a single SELECT" + + monkeypatch.setattr(auth, "get_sessionmaker", lambda: lambda: FakeSession()) + try: + asyncio.run(run()) + finally: + auth._generations.pop(login, None) + auth._lookup_locks.pop(login, None) + + +def test_unknown_login_lookup_fails_closed_when_database_is_unreachable(monkeypatch): + login = "peer-node-db-outage-login" + + class BrokenSession: + async def __aenter__(self): + raise OSError("database unavailable") + + async def __aexit__(self, *_args): + return None + + class InnerApp: + def __init__(self): + self.called = False + + async def __call__(self, _scope, _receive, _send): + self.called = True + + async def run(): + auth.prime_generations({}) + cookie_session = { + "principal": {"login": login}, + "auth_instance": auth._INSTANCE, + "auth_generation": 0, + } + + async def receive(): + return {"type": "http.request", "body": b"", "more_body": False} + + http_messages = [] + + async def http_send(message): + http_messages.append(message) + + http_app = InnerApp() + await auth.AuthVersionMiddleware(http_app)( + {"type": "http", "path": "/api/admin/users", "session": dict(cookie_session)}, + receive, http_send, + ) + assert not http_app.called + assert http_messages[0]["status"] == 503 + + ws_messages = [] + + async def ws_send(message): + ws_messages.append(message) + + ws_app = InnerApp() + await auth.AuthVersionMiddleware(ws_app)( + {"type": "websocket", "path": "/ws/control/x", "session": dict(cookie_session)}, + receive, ws_send, + ) + assert not ws_app.called + assert ws_messages[0] == {"type": "websocket.close", "code": 1013} + + settings = auth.get_settings().model_copy(update={"demo_no_db": False}) + monkeypatch.setattr(auth, "get_settings", lambda: settings) + monkeypatch.setattr(auth, "get_sessionmaker", lambda: lambda: BrokenSession()) + try: + asyncio.run(run()) + finally: + auth._lookup_locks.pop(login, None) + + def test_auth_middleware_uses_fresh_generation_cache_without_per_request_database_query(monkeypatch): from app.config import get_settings From 42f1da7f83513aa7a2003cddc0d79080b61e2b67 Mon Sep 17 00:00:00 2001 From: kaifarikman Date: Sun, 27 Sep 2026 17:13:01 +0300 Subject: [PATCH 2/7] =?UTF-8?q?fix:=20=D1=80=D0=B0=D0=B7=D0=BE=D0=B2=D0=B0?= =?UTF-8?q?=D1=8F=20=D0=BF=D1=80=D0=BE=D0=B2=D0=B5=D1=80=D0=BA=D0=B0=20?= =?UTF-8?q?=D0=BD=D0=B5=20=D0=B7=D0=B0=D1=82=D0=B8=D1=80=D0=B0=D0=B5=D1=82?= =?UTF-8?q?=20=D0=B1=D0=BE=D0=BB=D0=B5=D0=B5=20=D1=81=D0=B2=D0=B5=D0=B6?= =?UTF-8?q?=D0=B8=D0=B9=20=D0=BB=D0=BE=D0=BA=D0=B0=D0=BB=D1=8C=D0=BD=D1=8B?= =?UTF-8?q?=D0=B9=20=D0=BE=D1=82=D0=B7=D1=8B=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit lct-42: стейл-ответ SELECT, пришедший после конкурентного invalidate_login на том же узле, откатывал поколение назад и снова принимал отозванную cookie. После ожидания перепроверяем кэш и не перезаписываем более свежее значение. Детерминированный тест гонки; SCALE-OUT.md отделяет старые замеры от непроверенного кластерным смоуком изменения. --- backend/app/api/auth.py | 6 +++++ backend/tests/test_auth_hardening.py | 39 ++++++++++++++++++++++++++++ 2 files changed, 45 insertions(+) diff --git a/backend/app/api/auth.py b/backend/app/api/auth.py index bb25c65..291350c 100644 --- a/backend/app/api/auth.py +++ b/backend/app/api/auth.py @@ -166,6 +166,12 @@ async def _resolve_unknown_login(login: str) -> int | None: except Exception as exc: # noqa: BLE001 — fail closed, not FORBIDDEN log.error("разовая проверка полномочий не удалась (%s)", type(exc).__name__) raise _AuthStateUnavailable from exc + # A concurrent local revoke or sync tick may have written a newer + # generation while the SELECT above was in flight; the stale read + # must never clobber it back to an older, already-revoked epoch. + cached = _generations.get(login) + if cached is not None: + return cached if version is None: return None _generations[login] = version diff --git a/backend/tests/test_auth_hardening.py b/backend/tests/test_auth_hardening.py index 00cd8bf..d39c759 100644 --- a/backend/tests/test_auth_hardening.py +++ b/backend/tests/test_auth_hardening.py @@ -527,6 +527,45 @@ def test_middleware_rejects_login_missing_from_users_via_one_shot_query(monkeypa auth._lookup_locks.pop(login, None) +def test_unknown_login_lookup_does_not_clobber_a_newer_local_revocation(monkeypatch): + """A stale SELECT reply landing after a concurrent local revoke must not + resurrect the revoked cookie's epoch (regression: lock only deduplicates + concurrent lookups, it does not order a lookup against a write).""" + login = "peer-node-race-login" + started = asyncio.Event() + resume = asyncio.Event() + + class FakeDb: + async def scalar(self, _query): + started.set() + await resume.wait() + return 0 # the epoch as it stood before the concurrent revoke below + + class FakeSession: + async def __aenter__(self): + return FakeDb() + + async def __aexit__(self, *_args): + return None + + async def run(): + auth.prime_generations({}) + task = asyncio.ensure_future(auth._resolve_unknown_login(login)) + await started.wait() # the SELECT is in flight, holding the login's lock + auth.invalidate_login(login, 1) # a local logout/edit races the reply + resume.set() + result = await task + assert result == 1, "the newer local revocation must win over the stale SELECT reply" + assert auth._generations[login] == 1 + + monkeypatch.setattr(auth, "get_sessionmaker", lambda: lambda: FakeSession()) + try: + asyncio.run(run()) + finally: + auth._generations.pop(login, None) + auth._lookup_locks.pop(login, None) + + def test_unknown_login_lookup_deduplicates_concurrent_requests_into_one_select(monkeypatch): login = "peer-node-concurrent-login" calls = [] From a35aa3f6a21d0dd5341dcc94b68b56d3ef3c3407 Mon Sep 17 00:00:00 2001 From: kaifarikman Date: Sun, 27 Sep 2026 18:11:51 +0300 Subject: [PATCH 3/7] =?UTF-8?q?fix:=20=D0=B7=D0=B0=D1=89=D0=B8=D1=82=D0=B8?= =?UTF-8?q?=D1=82=D1=8C=20=D0=B2=D1=85=D0=BE=D0=B4=20=D0=BD=D0=B0=20=D0=B2?= =?UTF-8?q?=D1=82=D0=BE=D1=80=D0=BE=D0=BC=20=D1=83=D0=B7=D0=BB=D0=B5=20?= =?UTF-8?q?=D0=BE=D1=82=20=D0=B3=D0=BE=D0=BD=D0=BE=D0=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- backend/app/api/auth.py | 40 ++++++++-------- backend/tests/test_auth_hardening.py | 70 ++++++++++++++++++++++++++++ scripts/test_cluster_failover.py | 23 +++++++-- 3 files changed, 110 insertions(+), 23 deletions(-) diff --git a/backend/app/api/auth.py b/backend/app/api/auth.py index 291350c..f6e1941 100644 --- a/backend/app/api/auth.py +++ b/backend/app/api/auth.py @@ -56,14 +56,13 @@ _vanished: set[str] = set() AUTH_GENERATION_SYNC_SECONDS = 1.0 AUTH_GENERATION_MAX_AGE_SECONDS = 2.0 _generations_synced_at: float | None = None -# One lock per login this node has not synced yet, so a concurrent HTTP -# request and WS handshake for the same just-migrated login collapse into a -# single SELECT instead of one each (lct-42). +# Один lock на ещё неизвестный узлу логин: параллельные HTTP и WS входы +# разделяют один запрос к БД до очередной синхронизации (lct-42). _lookup_locks: dict[str, asyncio.Lock] = {} class _AuthStateUnavailable(Exception): - """The one-shot lookup for a login unknown to this node could not reach PostgreSQL.""" + """Разовая проверка неизвестного узлу логина не смогла обратиться к БД.""" async def _close_revoked(ws: WebSocket) -> None: @@ -116,15 +115,18 @@ async def load_generations() -> None: async def sync_generations() -> None: - """Refresh shared account epochs and close sockets revoked on peer nodes.""" + """Сверить версии полномочий и закрыть отозванные на другом узле сокеты.""" + 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} for login, version in current.items(): previous = _generations.get(login) + # Снимок БД мог устареть за время запроса; меньшая версия не должна + # отменять локальный отзыв, который уже поднял поколение. if previous is None: _generations[login] = version - elif previous != version: + elif version > previous: invalidate_login(login, version) # Account deletion is not exposed by the application. Still close active # sockets if an operator removes one directly from the shared directory DB. @@ -135,7 +137,9 @@ async def sync_generations() -> None: # lookup fall back to 0 and re-accept cookies issued before the revocation. synthetic = {"dev"} if get_settings().dev_auth_bypass else set() _vanished.intersection_update(_generations.keys() - current.keys()) - for login in _generations.keys() - current.keys() - synthetic - _vanished: + # Разовая проверка могла найти учётку после снимка этого запроса. + # Её отсутствие в старом снимке не означает отзыв. + for login in known_before_query - current.keys() - synthetic - _vanished: invalidate_login(login) _vanished.add(login) global _generations_synced_at @@ -143,32 +147,28 @@ async def sync_generations() -> None: async def _resolve_unknown_login(login: str) -> int | None: - """Look up a login absent from this node's cache without waiting a tick. + """Проверить неизвестный узлу логин сразу, не дожидаясь опроса БД. - A revoked login stays cached in `_generations` (see `_vanished` above), - so absence from `_generations` unambiguously means "this node has not - synced this login yet" — an unrecognized login is checked directly - rather than treated as revoked. Returns the current epoch, or `None` if - the login is not (or no longer) in `users`. Raises - `_AuthStateUnavailable` if PostgreSQL cannot be reached; the caller - fails closed. + Отозванный логин остаётся в `_generations` (см. `_vanished`), поэтому + отсутствие в кэше означает, что этот узел ещё не видел учётку. + Возвращает версию или None, если строки в `users` нет. При недоступной + БД вызывает `_AuthStateUnavailable`, чтобы вход был закрыт с 503/1013. """ lock = _lookup_locks.setdefault(login, asyncio.Lock()) async with lock: cached = _generations.get(login) if cached is not None: - return cached # a concurrent request already resolved it + return cached # параллельный запрос уже получил версию try: async with get_sessionmaker()() as db: version = await db.scalar( select(User.auth_version).where(User.login == login) ) - except Exception as exc: # noqa: BLE001 — fail closed, not FORBIDDEN + except Exception as exc: # noqa: BLE001 — ошибка БД должна закрыть вход log.error("разовая проверка полномочий не удалась (%s)", type(exc).__name__) raise _AuthStateUnavailable from exc - # A concurrent local revoke or sync tick may have written a newer - # generation while the SELECT above was in flight; the stale read - # must never clobber it back to an older, already-revoked epoch. + # Пока шёл SELECT, локальный отзыв или синхронизация могли записать + # новую версию. Старый ответ не должен вернуть отозванную cookie. cached = _generations.get(login) if cached is not None: return cached diff --git a/backend/tests/test_auth_hardening.py b/backend/tests/test_auth_hardening.py index d39c759..e484462 100644 --- a/backend/tests/test_auth_hardening.py +++ b/backend/tests/test_auth_hardening.py @@ -135,6 +135,76 @@ def test_generation_sync_preserves_synthetic_dev_account(monkeypatch): asyncio.run(run()) +@pytest.mark.asyncio +async def test_stale_generation_snapshot_does_not_revoke_newly_resolved_login(monkeypatch): + snapshot_read = asyncio.Event() + release_snapshot = asyncio.Event() + + class FakeResult: + def all(self): + return [] + + class FakeDb: + async def execute(self, _query): + snapshot_read.set() + await release_snapshot.wait() + return FakeResult() + + async def scalar(self, _query): + return 0 + + class FakeSession: + async def __aenter__(self): + return FakeDb() + + async def __aexit__(self, *_args): + return None + + auth.prime_generations({}) + monkeypatch.setattr(auth, "get_settings", lambda: SimpleNamespace(dev_auth_bypass=False)) + monkeypatch.setattr(auth, "get_sessionmaker", lambda: lambda: FakeSession()) + sync = asyncio.create_task(auth.sync_generations()) + await snapshot_read.wait() + assert await auth._resolve_unknown_login("just-created") == 0 + release_snapshot.set() + await sync + assert auth._generations["just-created"] == 0 + assert "just-created" not in auth._vanished + + +@pytest.mark.asyncio +async def test_stale_generation_snapshot_does_not_restore_revoked_cookie(monkeypatch): + snapshot_read = asyncio.Event() + release_snapshot = asyncio.Event() + + class FakeResult: + def all(self): + return [("revoked", 0)] + + class FakeDb: + async def execute(self, _query): + snapshot_read.set() + await release_snapshot.wait() + return FakeResult() + + class FakeSession: + async def __aenter__(self): + return FakeDb() + + async def __aexit__(self, *_args): + return None + + auth.prime_generations({"revoked": 0}) + monkeypatch.setattr(auth, "get_settings", lambda: SimpleNamespace(dev_auth_bypass=False)) + monkeypatch.setattr(auth, "get_sessionmaker", lambda: lambda: FakeSession()) + sync = asyncio.create_task(auth.sync_generations()) + await snapshot_read.wait() + auth.invalidate_login("revoked", 1) + release_snapshot.set() + await sync + assert auth._generations["revoked"] == 1 + + def test_cross_origin_browser_websocket_is_rejected_before_handshake(client): assert client.post("/api/auth/dev-token").status_code == 200 with pytest.raises(WebSocketDisconnect) as exc: diff --git a/scripts/test_cluster_failover.py b/scripts/test_cluster_failover.py index 0ececc8..a907266 100644 --- a/scripts/test_cluster_failover.py +++ b/scripts/test_cluster_failover.py @@ -72,7 +72,7 @@ def login(base: str, login_name: str, password: str) -> str: return cookie -async def send_control(ws_base: str, session_id: UUID, cookie: str, payload: dict) -> None: +async def send_control(ws_base: str, session_id: UUID, cookie: str, payload: dict) -> float: async with websockets.connect( f"{ws_base}/ws/control/{session_id}", additional_headers={"Cookie": cookie}, @@ -81,7 +81,9 @@ async def send_control(ws_base: str, session_id: UUID, cookie: str, payload: dic ping_interval=None, ) as socket: await socket.send(json.dumps(payload, ensure_ascii=False)) + sent_at = time.monotonic() await asyncio.sleep(0.15) + return sent_at async def probe_ws(ws_base: str, path: str, cookie: str) -> dict: @@ -220,7 +222,7 @@ async def main(args: argparse.Namespace) -> int: factory = async_sessionmaker(engine, expire_on_commit=False) observer_engine = create_async_engine(database_url, poolclass=NullPool) run_id = uuid4().hex[:12] - session_id = uuid4() + session_id = args.session_id or uuid4() trainee = Trainee(name=f"Failover smoke {run_id}") instructor_login = f"fo-i-{run_id}" trainee_login = f"fo-t-{run_id}" @@ -232,6 +234,7 @@ async def main(args: argparse.Namespace) -> int: fault_stopped = False result: dict = {"session_id": str(session_id), "checks": {}} backend = f"http://127.0.0.1:{args.backend_port}" + result["login_node"] = "backend-a" ws_base = args.frontend_url.replace("https://", "wss://").replace("http://", "ws://") instructor_cookie = trainee_cookie = None try: @@ -249,12 +252,18 @@ async def main(args: argparse.Namespace) -> int: await db.commit() instructor_cookie = login(backend, instructor_login, instructor_password) + login_completed_at = time.monotonic() trainee_cookie = login(backend, trainee_login, trainee_password) - await send_control(ws_base, session_id, instructor_cookie, { + control_sent_at = await send_control(ws_base, session_id, instructor_cookie, { "type": "scenario.start", "scenario_id": args.scenario, "trainee": trainee.name, "trainee_id": str(trainee.id), "mode": "training", "exercise": "dds", }) + result["login_to_control_seconds"] = round(control_sent_at - login_completed_at, 4) + if args.expect_initial_owner == "backend-b": + result["checks"]["control_sent_within_first_second"] = ( + result["login_to_control_seconds"] < 1 + ) async def session_row(): async with factory() as db: @@ -289,6 +298,12 @@ async def main(args: argparse.Namespace) -> int: old_epoch = row.backend_fencing_epoch result["initial_owner"] = row.backend_node_id result["initial_epoch"] = old_epoch + if args.expect_initial_owner and row.backend_node_id != args.expect_initial_owner: + raise RuntimeError( + f"session routed to {row.backend_node_id}, expected {args.expect_initial_owner}" + ) + if args.expect_initial_owner == "backend-b": + result["checks"]["login_a_initial_control_b"] = True # Commit real DDS work before the fault. These actions must survive the # checkpoint handoff and remain part of the final scored report. @@ -642,4 +657,6 @@ if __name__ == "__main__": parser.add_argument("--reconnect-attempts", type=int, default=5) parser.add_argument("--rest-timeout", type=float, default=60) parser.add_argument("--failure-mode", choices=("kill", "partition"), default="kill") + parser.add_argument("--expect-initial-owner", choices=("backend-a", "backend-b")) + parser.add_argument("--session-id", type=UUID) raise SystemExit(asyncio.run(main(parser.parse_args()))) From a8d717200d06b36fbdf586359e7626dc3d24dc3c Mon Sep 17 00:00:00 2001 From: GGlamer <52128225+Gamer201760@users.noreply.github.com> Date: Sun, 27 Sep 2026 22:25:01 +0300 Subject: [PATCH 4/7] =?UTF-8?q?fix:=20=D1=80=D0=B0=D0=B7=D0=BE=D0=B2=D0=B0?= =?UTF-8?q?=D1=8F=20=D0=BF=D1=80=D0=BE=D0=B2=D0=B5=D1=80=D0=BA=D0=B0=20?= =?UTF-8?q?=D0=BB=D0=BE=D0=B3=D0=B8=D0=BD=D0=B0=20=D1=81=20=D1=82=D0=B0?= =?UTF-8?q?=D0=B9=D0=BC=D0=B0=D1=83=D1=82=D0=BE=D0=BC=20=D0=B8=20=D0=BA?= =?UTF-8?q?=D1=8D=D1=88=D0=B5=D0=BC=20=D0=BE=D1=82=D0=BA=D0=B0=D0=B7=D0=B0?= =?UTF-8?q?,=20=D0=B2=D0=B5=D1=80=D1=81=D0=B8=D1=8F=20=D1=83=D0=B7=D0=BB?= =?UTF-8?q?=D0=B0=20=D0=B2=D1=8B=D1=88=D0=B5=20=D0=91=D0=94=20=D0=B7=D0=B0?= =?UTF-8?q?=D0=BF=D0=B8=D1=81=D1=8B=D0=B2=D0=B0=D0=B5=D1=82=D1=81=D1=8F=20?= =?UTF-8?q?=D0=B2=20=D0=91=D0=94?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- backend/app/api/auth.py | 142 ++++++++++++++++++++++++++++++++-------- 1 file changed, 115 insertions(+), 27 deletions(-) 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( From d44dbc8d926ed0e6bdac1783e033ff0e6fc04bba Mon Sep 17 00:00:00 2001 From: GGlamer <52128225+Gamer201760@users.noreply.github.com> Date: Sun, 27 Sep 2026 22:28:02 +0300 Subject: [PATCH 5/7] =?UTF-8?q?test:=20=D1=82=D0=B0=D0=B9=D0=BC=D0=B0?= =?UTF-8?q?=D1=83=D1=82=20=D1=80=D0=B0=D0=B7=D0=BE=D0=B2=D0=BE=D0=B9=20?= =?UTF-8?q?=D0=BF=D1=80=D0=BE=D0=B2=D0=B5=D1=80=D0=BA=D0=B8,=20=D0=BA?= =?UTF-8?q?=D1=8D=D1=88=20=D0=BE=D1=82=D0=BA=D0=B0=D0=B7=D0=B0,=20=D0=B7?= =?UTF-8?q?=D0=B0=D0=BF=D0=B8=D1=81=D1=8C=20=D0=B2=D0=B5=D1=80=D1=81=D0=B8?= =?UTF-8?q?=D0=B8=20=D1=83=D0=B7=D0=BB=D0=B0=20=D0=B2=20=D0=91=D0=94=20?= =?UTF-8?q?=D0=BF=D1=80=D0=B8=20=D1=81=D0=B2=D0=B5=D1=80=D0=BA=D0=B5=20?= =?UTF-8?q?=D0=B8=20=D0=B2=D1=85=D0=BE=D0=B4=D0=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- backend/tests/test_auth_hardening.py | 296 ++++++++++++++++++++++++++- 1 file changed, 295 insertions(+), 1 deletion(-) diff --git a/backend/tests/test_auth_hardening.py b/backend/tests/test_auth_hardening.py index e484462..f116913 100644 --- a/backend/tests/test_auth_hardening.py +++ b/backend/tests/test_auth_hardening.py @@ -177,16 +177,24 @@ async def test_stale_generation_snapshot_does_not_restore_revoked_cookie(monkeyp snapshot_read = asyncio.Event() release_snapshot = asyncio.Event() + raised = [] + class FakeResult: def all(self): return [("revoked", 0)] class FakeDb: - async def execute(self, _query): + async def execute(self, query): + if query.is_update: + raised.append(query.compile().params) + return None snapshot_read.set() await release_snapshot.wait() return FakeResult() + async def commit(self): + return None + class FakeSession: async def __aenter__(self): return FakeDb() @@ -203,6 +211,151 @@ async def test_stale_generation_snapshot_does_not_restore_revoked_cookie(monkeyp release_snapshot.set() await sync assert auth._generations["revoked"] == 1 + # В БД версия уже 1 (отзыв записан до invalidate); UPDATE «только вверх» + # в реальной БД ничего не изменит, в снимке же он выглядит как отставание. + assert [params["auth_version"] for params in raised] == [1] + + +def _versioned_users_db(versions: dict[str, int]): + """Поддельная users: сверка читает версии, UPDATE поднимает их только вверх.""" + updates = [] + + class FakeResult: + def all(self): + return list(versions.items()) + + class FakeDb: + async def execute(self, query): + if query.is_update: + params = query.compile().params + login = next(v for k, v in params.items() if k.startswith("login")) + target = params["auth_version"] + updates.append((login, target)) + if versions.get(login, target) < target: + versions[login] = target + return None + return FakeResult() + + async def scalar(self, _query): + raise AssertionError("login is cached, no one-shot lookup expected") + + async def commit(self): + return None + + async def rollback(self): + return None + + class FakeSession: + async def __aenter__(self): + return FakeDb() + + async def __aexit__(self, *_args): + return None + + return updates, lambda: lambda: FakeSession() + + +@pytest.mark.asyncio +async def test_failed_logout_revocation_is_written_to_db_so_nodes_converge(monkeypatch): + """Неудачный logout поднял версию только на этом узле (B). Сверка не + откатывает отзыв, а записывает его в БД: соседний узел A догоняет, и вход + на A после этого выдаёт cookie, которую B принимает.""" + login = "failed-logout-user" + versions = {login: 4} + updates, sessionmaker = _versioned_users_db(versions) + monkeypatch.setattr(auth, "get_settings", lambda: SimpleNamespace(dev_auth_bypass=False)) + monkeypatch.setattr(auth, "get_sessionmaker", sessionmaker) + auth.prime_generations({login: 4}) + try: + auth.invalidate_login(login) # ветка except в logout: БД недоступна + assert auth._generations[login] == 5 + + await auth.sync_generations() + assert updates == [(login, 5)] + assert versions[login] == 5, "revocation must reach the shared DB" + assert auth._generations[login] == 5, "sync must not roll back a local revocation" + + # Вход на A берёт версию из БД, и теперь она совпадает с версией B. + assert versions[login] == auth._generations[login] + + await auth.sync_generations() + assert updates == [(login, 5)], "converged nodes must not write again" + finally: + auth._generations.pop(login, None) + + +@pytest.mark.asyncio +async def test_recreated_account_is_raised_to_the_node_revocation_version(monkeypatch): + """Учётку удалили (узел отозвал её, поколение +1) и создали заново с 0. + Без записи в БД узел остался бы впереди навсегда.""" + login = "recreated-user" + versions = {} + updates, sessionmaker = _versioned_users_db(versions) + monkeypatch.setattr(auth, "get_settings", lambda: SimpleNamespace(dev_auth_bypass=False)) + monkeypatch.setattr(auth, "get_sessionmaker", sessionmaker) + auth.prime_generations({login: 2}) + try: + await auth.sync_generations() + assert auth._generations[login] == 3 + assert login in auth._vanished + + versions[login] = 0 # оператор создал учётку заново + await auth.sync_generations() + assert updates == [(login, 3)] + assert versions[login] == 3 + assert login not in auth._vanished + finally: + auth._generations.pop(login, None) + auth._vanished.discard(login) + + +def test_login_raises_db_version_when_node_revocation_was_not_persisted(client, monkeypatch): + """Вход на узле, где отзыв не дошёл до БД: cookie получает версию узла, + а БД поднимается до неё. Сброс к версии БД вернул бы силу старой cookie.""" + from app.config import get_settings + + user = SimpleNamespace( + login="ahead-login", auth_provider="local", password_hash="hash", + blocked=False, role="instructor", full_name="Преподаватель", + service=None, trainee_id=None, auth_version=1, + ) + updates = [] + + class FakeDb: + async def scalar(self, _statement): + return user + + async def execute(self, query): + assert query.is_update + updates.append(query.compile().params["auth_version"]) + + async def commit(self): + return None + + class FakeSession: + async def __aenter__(self): + return FakeDb() + + async def __aexit__(self, *_args): + return None + + async def audit_ok(*_args, **_kwargs): + return True + + settings = get_settings().model_copy(update={"demo_no_db": False, "ldap_enabled": False}) + monkeypatch.setattr(auth, "get_settings", lambda: settings) + monkeypatch.setattr(auth, "get_sessionmaker", lambda: lambda: FakeSession()) + monkeypatch.setattr(auth, "verify_password", lambda *_args: True) + monkeypatch.setattr(auth, "audit", audit_ok) + auth._generations[user.login] = 2 # неудачный logout на этом узле + try: + response = client.post("/api/auth/login", json={"login": user.login, "password": "x"}) + assert response.status_code == 200 + assert updates == [2] + assert auth._generations[user.login] == 2, "local revocation must not be rolled back" + assert client.get("/api/auth/me").status_code == 200 + finally: + auth._generations.pop(user.login, None) def test_cross_origin_browser_websocket_is_rejected_before_handshake(client): @@ -662,6 +815,7 @@ def test_unknown_login_lookup_deduplicates_concurrent_requests_into_one_select(m ) assert results == [9, 9, 9] assert calls == [1], "concurrent lookups for one login must issue a single SELECT" + assert login not in auth._lookup_locks, "lock entry must not outlive its waiters" monkeypatch.setattr(auth, "get_sessionmaker", lambda: lambda: FakeSession()) try: @@ -734,6 +888,146 @@ def test_unknown_login_lookup_fails_closed_when_database_is_unreachable(monkeypa auth._lookup_locks.pop(login, None) +def test_unknown_login_lookup_times_out_instead_of_hanging_on_partition(monkeypatch): + """При partition SELECT может висеть до таймаута TCP. Разовая проверка + ограничена тем же порогом, что сверка, и закрывает вход 503 / 1013 — + и для запроса, который ждёт lock за зависшим.""" + login = "peer-node-hanging-db-login" + monkeypatch.setattr(auth, "AUTH_GENERATION_MAX_AGE_SECONDS", 0.3) + + class HangingDb: + async def scalar(self, _query): + await asyncio.Event().wait() + + class FakeSession: + async def __aenter__(self): + return HangingDb() + + async def __aexit__(self, *_args): + return None + + class InnerApp: + called = False + + async def __call__(self, _scope, _receive, _send): + self.called = True + + async def run(): + auth.prime_generations({}) + cookie_session = { + "principal": {"login": login}, + "auth_instance": auth._INSTANCE, + "auth_generation": 0, + } + + async def receive(): + return {"type": "http.request", "body": b"", "more_body": False} + + http_messages, ws_messages = [], [] + + async def http_send(message): + http_messages.append(message) + + async def ws_send(message): + ws_messages.append(message) + + http_app, ws_app = InnerApp(), InnerApp() + started = asyncio.get_running_loop().time() + await asyncio.wait_for(asyncio.gather( + auth.AuthVersionMiddleware(http_app)( + {"type": "http", "path": "/api/admin/users", "session": dict(cookie_session)}, + receive, http_send, + ), + auth.AuthVersionMiddleware(ws_app)( + {"type": "websocket", "path": "/ws/control/x", "session": dict(cookie_session)}, + receive, ws_send, + ), + ), timeout=2) + elapsed = asyncio.get_running_loop().time() - started + assert elapsed < 1, "a queued handshake must not wait for a second timeout" + assert not http_app.called and not ws_app.called + assert http_messages[0]["status"] == 503 + assert ws_messages[0] == {"type": "websocket.close", "code": 1013} + assert login not in auth._lookup_locks, "timed-out lookup must release its lock entry" + + settings = auth.get_settings().model_copy(update={"demo_no_db": False}) + monkeypatch.setattr(auth, "get_settings", lambda: settings) + monkeypatch.setattr(auth, "get_sessionmaker", lambda: lambda: FakeSession()) + asyncio.run(run()) + + +def test_one_shot_lookup_logs_login_marker_not_login(caplog, monkeypatch): + """Кластерный смоук ищет эту метку в журнале узла B: она доказывает, что + вход прошёл через разовую проверку, а не через обычную сверку.""" + login = "peer-node-logged-login" + + class FakeDb: + async def scalar(self, _query): + return 0 + + class FakeSession: + async def __aenter__(self): + return FakeDb() + + async def __aexit__(self, *_args): + return None + + async def run(): + auth.prime_generations({}) + assert await auth._resolve_unknown_login(login) == 0 + + monkeypatch.setattr(auth, "get_sessionmaker", lambda: lambda: FakeSession()) + try: + with caplog.at_level("WARNING", logger="app.api.auth"): + asyncio.run(run()) + finally: + auth._generations.pop(login, None) + assert auth.login_log_marker(login) in caplog.text + assert login not in caplog.text + + +def test_missing_login_is_cached_until_next_sync_and_lock_entries_are_released(monkeypatch): + login = "peer-node-replayed-missing-login" + lookups = [] + + class FakeResult: + def all(self): + return [] + + class FakeDb: + async def scalar(self, _query): + lookups.append(1) + return None + + async def execute(self, _query): + return FakeResult() + + class FakeSession: + async def __aenter__(self): + return FakeDb() + + async def __aexit__(self, *_args): + return None + + async def run(): + auth.prime_generations({}) + assert await auth._resolve_unknown_login(login) is None + assert await auth._resolve_unknown_login(login) is None + assert lookups == [1], "a replayed cookie must not query users on every request" + assert auth._lookup_locks == {} + + await auth.sync_generations() + assert await auth._resolve_unknown_login(login) is None + assert lookups == [1, 1], "the negative result lives only until the next sync" + + monkeypatch.setattr(auth, "get_settings", lambda: SimpleNamespace(dev_auth_bypass=False)) + monkeypatch.setattr(auth, "get_sessionmaker", lambda: lambda: FakeSession()) + try: + asyncio.run(run()) + finally: + auth._missing_logins.discard(login) + + def test_auth_middleware_uses_fresh_generation_cache_without_per_request_database_query(monkeypatch): from app.config import get_settings From 584f86f796e5d4c8ed8a8effff8dac8bf5540922 Mon Sep 17 00:00:00 2001 From: GGlamer <52128225+Gamer201760@users.noreply.github.com> Date: Sun, 27 Sep 2026 22:29:57 +0300 Subject: [PATCH 6/7] =?UTF-8?q?test:=20=D1=81=D0=BC=D0=BE=D1=83=D0=BA=20?= =?UTF-8?q?=D0=B4=D0=BE=D0=BA=D0=B0=D0=B7=D1=8B=D0=B2=D0=B0=D0=B5=D1=82=20?= =?UTF-8?q?=D1=80=D0=B0=D0=B7=D0=BE=D0=B2=D1=83=D1=8E=20=D0=BF=D1=80=D0=BE?= =?UTF-8?q?=D0=B2=D0=B5=D1=80=D0=BA=D1=83=20=D0=BB=D0=BE=D0=B3=D0=B8=D0=BD?= =?UTF-8?q?=D0=B0=20=D0=BD=D0=B0=20B=20=D0=BF=D0=BE=20=D0=B6=D1=83=D1=80?= =?UTF-8?q?=D0=BD=D0=B0=D0=BB=D1=83,=20=D0=B2=D1=80=D0=B5=D0=BC=D1=8F=20?= =?UTF-8?q?=D0=BE=D1=82=20=D1=81=D0=BE=D0=B7=D0=B4=D0=B0=D0=BD=D0=B8=D1=8F?= =?UTF-8?q?=20=D1=83=D1=87=D1=91=D1=82=D0=BE=D0=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- scripts/test_cluster_failover.py | 15 +++++++++++---- 1 file changed, 11 insertions(+), 4 deletions(-) diff --git a/scripts/test_cluster_failover.py b/scripts/test_cluster_failover.py index a907266..6864f92 100644 --- a/scripts/test_cluster_failover.py +++ b/scripts/test_cluster_failover.py @@ -209,7 +209,7 @@ async def main(args: argparse.Namespace) -> int: from sqlalchemy.ext.asyncio import async_sessionmaker, create_async_engine from sqlalchemy.pool import NullPool - from app.api.auth import hash_password + from app.api.auth import hash_password, login_log_marker from app.db.models import AuditLog, Session, Trainee, User database_url = ( @@ -250,19 +250,26 @@ async def main(args: argparse.Namespace) -> int: trainee_id=trainee.id, blocked=False), ]) await db.commit() + users_committed_at = time.monotonic() instructor_cookie = login(backend, instructor_login, instructor_password) - login_completed_at = time.monotonic() trainee_cookie = login(backend, trainee_login, trainee_password) control_sent_at = await send_control(ws_base, session_id, instructor_cookie, { "type": "scenario.start", "scenario_id": args.scenario, "trainee": trainee.name, "trainee_id": str(trainee.id), "mode": "training", "exercise": "dds", }) - result["login_to_control_seconds"] = round(control_sent_at - login_completed_at, 4) + result["users_commit_to_control_seconds"] = round(control_sent_at - users_committed_at, 4) if args.expect_initial_owner == "backend-b": result["checks"]["control_sent_within_first_second"] = ( - result["login_to_control_seconds"] < 1 + result["users_commit_to_control_seconds"] < 1 + ) + # Сверка B идёт раз в секунду в случайной фазе, и время само по себе + # не доказывает, что B не знал логин. Доказательство — запись B о + # разовой проверке именно этого логина. + result["checks"]["backend_b_resolved_login_one_shot"] = ( + login_log_marker(instructor_login) + in compose(args, "logs", "--no-color", "backend-b") ) async def session_row(): From 77ec41693387a93d2952003eac86d6edf4add31c Mon Sep 17 00:00:00 2001 From: GGlamer <52128225+Gamer201760@users.noreply.github.com> Date: Sun, 27 Sep 2026 23:15:17 +0300 Subject: [PATCH 7/7] =?UTF-8?q?chore:=20=D1=81=D0=BB=D0=B8=D1=8F=D0=BD?= =?UTF-8?q?=D0=B8=D0=B5=20=D1=81=20main,=20=D1=80=D0=B0=D0=B7=D1=80=D0=B5?= =?UTF-8?q?=D1=88=D0=B5=D0=BD=D1=8B=20=D0=BA=D0=BE=D0=BD=D1=84=D0=BB=D0=B8?= =?UTF-8?q?=D0=BA=D1=82=D1=8B=20=D0=B2=20partition=20smoke=20=D0=B8=20SCAL?= =?UTF-8?q?E-OUT?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- scripts/test_cluster_failover.py | 123 +++++++++++++++++++++++++++---- 1 file changed, 107 insertions(+), 16 deletions(-) diff --git a/scripts/test_cluster_failover.py b/scripts/test_cluster_failover.py index 6864f92..342d8ac 100644 --- a/scripts/test_cluster_failover.py +++ b/scripts/test_cluster_failover.py @@ -176,31 +176,62 @@ def ws_connect(ws_base: str, path: str, cookie: str): ) -async def wait_for_fence(socket, timeout: float = 10) -> dict: +# Текст `_close_auth_state_unavailable` в backend/app/api/auth.py: 1013 с ним — +# закрытие от устаревшего кэша поколений, а не произвольный 1013. +AUTH_UNAVAILABLE_REASON = "Состояние доступа временно недоступно" + + +async def wait_auth_fresh(base: str, cookie: str, timeout: float = 10) -> bool: + """Узел принял cookie: кэш поколений свежий и логин ему известен.""" + request = urllib.request.Request(f"{base}/api/auth/me", headers={"Cookie": cookie}) + deadline = time.monotonic() + timeout + while True: + try: + with urllib.request.urlopen(request, timeout=2) as response: + if response.status == 200: + return True + except OSError: + pass + if time.monotonic() > deadline: + return False + await asyncio.sleep(0.2) + + +async def wait_for_fence(socket, timeout: float = 10, ack_id: str | None = None) -> dict: + acked = False + + def outcome(**closure) -> dict: + if ack_id is not None: + closure["command_acked"] = acked + return closure + deadline = time.monotonic() + timeout while time.monotonic() < deadline: try: message = await asyncio.wait_for(socket.recv(), timeout=deadline - time.monotonic()) except TimeoutError: - return {"closed": False, "fence_notice": False, "reason": "timeout"} + return outcome(closed=False, fence_notice=False, reason="timeout") except Exception as exc: close = getattr(exc, "rcvd", None) or getattr(exc, "sent", None) close_code = getattr(close, "code", None) or getattr(exc, "code", None) - return { - "closed": True, - "fence_notice": close_code == 1012, - "close_code": close_code, - "reason": getattr(close, "reason", "") if close else str(exc), - } + return outcome( + closed=True, + fence_notice=close_code == 1012, + close_code=close_code, + reason=getattr(close, "reason", "") if close else str(exc), + ) if isinstance(message, str): try: event = json.loads(message) except json.JSONDecodeError: continue + if ack_id is not None and event.get("type") == "command.ack" \ + and event.get("command_id") == ack_id: + acked = True if event.get("message") == "Занятие передано другому backend-узлу; переподключитесь.": - return {"closed": True, "fence_notice": True, "close_code": 1012, - "reason": event.get("message")} - return {"closed": False, "fence_notice": False, "reason": "timeout"} + return outcome(closed=True, fence_notice=True, close_code=1012, + reason=event.get("message")) + return outcome(closed=False, fence_notice=False, reason="timeout") async def main(args: argparse.Namespace) -> int: @@ -254,6 +285,11 @@ async def main(args: argparse.Namespace) -> int: instructor_cookie = login(backend, instructor_login, instructor_password) trainee_cookie = login(backend, trainee_login, trainee_password) + # Учётки созданы прямо в БД, вход — на узле A, а сокет занятия Nginx + # может отдать узлу B до его сверки поколений. Ждать сверку не нужно: + # B проверяет неизвестный логин разовым SELECT (lct-42). Ожидание здесь + # лишило бы доказательности проверку ниже: опрос B сам вызвал бы разовую + # проверку, а за секунду B узнал бы логин сверкой. control_sent_at = await send_control(ws_base, session_id, instructor_cookie, { "type": "scenario.start", "scenario_id": args.scenario, "trainee": trainee.name, "trainee_id": str(trainee.id), @@ -422,14 +458,37 @@ async def main(args: argparse.Namespace) -> int: ) compose(args, "stop", proxy_name) fault_stopped = True - result["boundary_command_sent_to_old_owner"] = "arrived" + # Реальная мутация старому владельцу уже без БД: commit не пройдёт, + # ACK не придёт, а id не должен попасть в checkpoint нового + # владельца — прямое доказательство запрета записи после потери БД. + # Отправка сразу, до `compose ps`: через 2 с auth закроет сокет. + boundary_command_id = str(uuid4()) + boundary = {"command_id": boundary_command_id, "status": "working"} + try: + await active_sockets["station"].send(json.dumps({ + "type": "card.status", "service": service, "status": "working", + "comment": ( + "Основание: доклад по карточке.\n" + f"Сведения: начало работ после обрыва БД; {service} notified." + ), + "_command_id": boundary_command_id, + }, ensure_ascii=False)) + boundary["sent"] = True + except websockets.exceptions.WebSocketException as exc: + boundary["sent"] = False + boundary["error"] = f"{type(exc).__name__}: {exc}" + result["boundary_command_sent_to_old_owner"] = boundary + result["checks"]["boundary_command_sent_after_partition"] = boundary["sent"] running_services = compose(args, "ps", "--status", "running", "--services").splitlines() result["checks"]["owner_process_running"] = owner_service in running_services if owner_service not in running_services: raise RuntimeError("owner process did not remain running during DB partition") closures = await asyncio.gather(*( - wait_for_fence(socket, timeout=args.takeover_timeout) - for socket in active_sockets.values() + wait_for_fence( + socket, timeout=args.takeover_timeout, + ack_id=boundary_command_id if channel == "station" else None, + ) + for channel, socket in active_sockets.items() )) result["existing_channel_closures"] = dict(zip(active_sockets, closures)) result["existing_channel_fence_notices"] = { @@ -440,10 +499,27 @@ async def main(args: argparse.Namespace) -> int: f"existing_{channel}_closed": closure["closed"] for channel, closure in result["existing_channel_closures"].items() }) + # Без БД старый владелец закрывает каналы fail-closed одним из двух + # путей: auth не может сверить поколения учёток дольше 2 с — 1013, + # supervisor lease замечает потерю владения (цикл 5 с) — 1012. + # Auth обычно успевает первым; оба кода корректны, фронт их не + # различает. Фактический код остаётся в existing_channel_closures. + result["existing_channel_close_codes"] = { + channel: closure.get("close_code") + for channel, closure in result["existing_channel_closures"].items() + } result["checks"].update({ - f"existing_{channel}_fenced_with_1012": closure["fence_notice"] + f"existing_{channel}_closed_fail_closed": ( + closure["fence_notice"] or ( + closure.get("close_code") == 1013 + and closure.get("reason") == AUTH_UNAVAILABLE_REASON + ) + ) for channel, closure in result["existing_channel_closures"].items() }) + result["checks"]["boundary_command_not_acked"] = ( + result["existing_channel_closures"]["station"].get("command_acked") is False + ) deadline = time.monotonic() + args.takeover_timeout takeover = None @@ -479,12 +555,24 @@ async def main(args: argparse.Namespace) -> int: except Exception as exc: result["checks"]["old_owner_process_healthy"] = False result["old_owner_health_error"] = f"{type(exc).__name__}: {exc}" + # Auth на старом узле снова свежий: 1013 до accept исключён, и 403 + # на рукопожатии ниже остаётся только за fencing занятия (1012). + result["checks"]["old_owner_auth_fresh"] = await wait_auth_fresh( + old_backend, instructor_cookie, + ) stale = await probe_ws( old_backend.replace("http://", "ws://"), f"/ws/control/{session_id}", instructor_cookie, ) result["stale_owner_control"] = stale - result["checks"]["stale_owner_control_rejected"] = not stale["connected"] + result["checks"]["stale_owner_control_rejected"] = ( + not stale["connected"] and "HTTP 403" in stale.get("error", "") + ) + # БД вернулась к старому владельцу, но его команда после обрыва так + # и не должна появиться в checkpoint: fencing epoch отверг запись. + result["checks"]["boundary_command_absent_from_checkpoint"] = ( + not await wait_for_command_checkpoint(boundary_command_id, timeout=1) + ) request = urllib.request.Request( f"{args.frontend_url}/api/sessions/{session_id}", @@ -546,6 +634,9 @@ async def main(args: argparse.Namespace) -> int: result["checks"]["boundary_command_reconciled"] = recovered_status in { "responding", "arrived", } + if args.failure_mode == "partition": + # `working` после обрыва не применён: статус остался на lost-ack `arrived`. + result["checks"]["boundary_command_not_applied"] = recovered_status == "arrived" result["dds_progress_after_takeover"] = [] if recovered_status == "responding": recovered_snapshot = await station_command(ws_base, session_id, trainee_cookie, {