diff --git a/backend/app/core/events.py b/backend/app/core/events.py new file mode 100644 index 0000000..e05f149 --- /dev/null +++ b/backend/app/core/events.py @@ -0,0 +1,56 @@ +"""In-memory SSE-шина (pub/sub) для пушей «данные изменились». + +Рассчитана на один процесс (uvicorn `--workers 1`). При переходе на несколько воркеров +шину нужно вынести во внешний брокер (например, Redis pub/sub). + +Публикация вызывается из СИНХРОННЫХ роутеров (FastAPI выполняет их в threadpool), а очереди +подписчиков живут в event-loop'е — поэтому публикация перекидывается в loop через +`call_soon_threadsafe`, а весь доступ к подпискам происходит на потоке loop'а. +""" +from __future__ import annotations + +import asyncio +from collections import defaultdict +from typing import Any, Iterable + + +class EventHub: + def __init__(self) -> None: + self._subs: dict[int, set[asyncio.Queue]] = defaultdict(set) + self._loop: asyncio.AbstractEventLoop | None = None + + def bind_loop(self, loop: asyncio.AbstractEventLoop) -> None: + self._loop = loop + + def subscribe(self, user_id: int) -> asyncio.Queue: + queue: asyncio.Queue = asyncio.Queue(maxsize=100) + self._subs[user_id].add(queue) + return queue + + def unsubscribe(self, user_id: int, queue: asyncio.Queue) -> None: + subs = self._subs.get(user_id) + if subs is not None: + subs.discard(queue) + if not subs: + self._subs.pop(user_id, None) + + def publish(self, user_ids: Iterable[int], event: dict[str, Any]) -> None: + """Доставить событие подписчикам с указанными user_id (потокобезопасно).""" + loop = self._loop + if loop is None: + return + ids = [uid for uid in user_ids if uid is not None] + if ids: + loop.call_soon_threadsafe(self._fanout, ids, event) + + def _fanout(self, ids: list[int], event: dict[str, Any]) -> None: + # Выполняется на потоке loop'а — доступ к _subs безопасен. + for uid in ids: + for queue in list(self._subs.get(uid, ())): + try: + queue.put_nowait(event) + except asyncio.QueueFull: + pass # медленный клиент — пропускаем (догонит при reconnect/focus) + + +hub = EventHub() diff --git a/backend/app/main.py b/backend/app/main.py index 0f33473..bd9484f 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -10,7 +10,6 @@ from fastapi import FastAPI, Request from fastapi.exceptions import RequestValidationError from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import FileResponse, JSONResponse -from starlette.middleware.base import BaseHTTPMiddleware from app.core import security from app.core.config import settings @@ -19,6 +18,7 @@ from app.routers import ( achievements, admin, auth, + events, groups, invitations, matches, @@ -33,36 +33,53 @@ _STATIC_DIR = Path(os.getenv("STATIC_DIR", str(Path(__file__).resolve().parent.p _UNSAFE_METHODS = {"POST", "PUT", "PATCH", "DELETE"} -class CSRFMiddleware(BaseHTTPMiddleware): - """Double-submit CSRF: для аутентифицированных мутаций на /api требуем - совпадения заголовка X-CSRF-Token и cookie csrf_token.""" +class CSRFMiddleware: + """Double-submit CSRF на чистом ASGI: для аутентифицированных мутаций на /api требуем + совпадения заголовка X-CSRF-Token и cookie csrf_token. - async def dispatch(self, request: Request, call_next): # noqa: ANN001 - path = request.url.path - if request.method in _UNSAFE_METHODS and path.startswith("/api"): - has_session = ( - security.USER_COOKIE in request.cookies - or security.ADMIN_COOKIE in request.cookies - ) - if has_session: - cookie_token = request.cookies.get(security.CSRF_COOKIE) - header_token = request.headers.get(security.CSRF_HEADER) - if not cookie_token or cookie_token != header_token: - return JSONResponse( - status_code=403, - content={ - "error": { - "code": "CSRF_FAILED", - "message": "Неверный или отсутствующий CSRF-токен.", - "details": None, - } - }, - ) - return await call_next(request) + Намеренно НЕ на BaseHTTPMiddleware: тот буферизует потоковые ответы и ломает SSE + (/api/events). Чистый ASGI пропускает стримы насквозь, вмешиваясь только при отказе CSRF. + """ + + def __init__(self, app) -> None: # noqa: ANN001 + self.app = app + + async def __call__(self, scope, receive, send): # noqa: ANN001 + if scope["type"] == "http": + request = Request(scope) + if request.method in _UNSAFE_METHODS and request.url.path.startswith("/api"): + has_session = ( + security.USER_COOKIE in request.cookies + or security.ADMIN_COOKIE in request.cookies + ) + if has_session: + cookie_token = request.cookies.get(security.CSRF_COOKIE) + header_token = request.headers.get(security.CSRF_HEADER) + if not cookie_token or cookie_token != header_token: + response = JSONResponse( + status_code=403, + content={ + "error": { + "code": "CSRF_FAILED", + "message": "Неверный или отсутствующий CSRF-токен.", + "details": None, + } + }, + ) + await response(scope, receive, send) + return + await self.app(scope, receive, send) @asynccontextmanager async def _lifespan(_app: FastAPI): + # SSE-шина публикует из sync-роутеров в этот event-loop — сохраняем ссылку (все окружения). + import asyncio + + from app.core.events import hub + + hub.bind_loop(asyncio.get_running_loop()) + # В DEV приложение само подтягивает справочники и админа из .env при старте # (в test/prod это делает entrypoint.sh; в pytest отключено FS_STARTUP_BOOTSTRAP=0). if settings.is_development and os.getenv("FS_STARTUP_BOOTSTRAP", "1") != "0": @@ -118,7 +135,7 @@ def create_app() -> FastAPI: # API-роутеры под /api. api_routers = [auth.router, users.router, groups.router, invitations.router, matches.router, reference.router, stats.router, achievements.router, - admin.router] + events.router, admin.router] for r in api_routers: app.include_router(r, prefix="/api") diff --git a/backend/app/routers/admin.py b/backend/app/routers/admin.py index dd5e639..09d2087 100644 --- a/backend/app/routers/admin.py +++ b/backend/app/routers/admin.py @@ -20,6 +20,7 @@ from app.services import ( audit_service, faction_service, match_service, + notify, user_service, ) @@ -229,6 +230,7 @@ def update_match( win_reason=body.win_reason, win_reason_set=("win_reason" in body.model_fields_set), participants=participants, + expected_version=body.expected_version, ) audit_service.record( session, @@ -239,6 +241,7 @@ def update_match( ip=request.client.host if request.client else None, ) session.commit() + notify.match_changed(session, match) return build_match_read(session, match, can_modify=True) @@ -288,12 +291,14 @@ def delete_match( session: Session = Depends(get_session), admin: User = Depends(get_current_admin), ) -> s.OkResponse: + group_id = match_service.get_match(session, match_id).group_id # для уведомления admin_service.delete_match(session, match_id) audit_service.record( session, actor_id=admin.id, action="delete", entity_type="match", entity_id=match_id, ip=request.client.host if request.client else None, ) session.commit() + notify.match_removed(session, match_id, group_id) return s.OkResponse() @@ -329,6 +334,7 @@ def admin_add_attachment( att = attachment_service.add_photo( session, match, admin, content, ext, user_service.avatar_media_type(ext) ) + notify.match_changed(session, match) return attachment_read(att, f"/api/admin/matches/{match_id}") @@ -341,6 +347,7 @@ def admin_delete_attachment( ) -> s.OkResponse: match = match_service.get_match(session, match_id) attachment_service.delete(session, match, attachment_id) + notify.match_changed(session, match) return s.OkResponse() diff --git a/backend/app/routers/events.py b/backend/app/routers/events.py new file mode 100644 index 0000000..e264f5b --- /dev/null +++ b/backend/app/routers/events.py @@ -0,0 +1,45 @@ +"""SSE: поток событий «данные изменились» для текущего пользователя. + +Клиент (EventSource) держит одно соединение и по событию точечно перезапрашивает данные +(push-to-invalidate). Мутации остаются обычным REST. +""" +from __future__ import annotations + +import asyncio +import json + +from fastapi import APIRouter, Depends, Request +from fastapi.responses import StreamingResponse + +from app.auth.deps import get_current_user +from app.core.events import hub +from app.models import User + +router = APIRouter(tags=["events"]) + +_HEARTBEAT_SECONDS = 25 + + +@router.get("/events") +async def events(request: Request, user: User = Depends(get_current_user)) -> StreamingResponse: + queue = hub.subscribe(user.id) # type: ignore[arg-type] + + async def stream(): + try: + yield ": connected\n\n" + while True: + if await request.is_disconnected(): + break + try: + event = await asyncio.wait_for(queue.get(), timeout=_HEARTBEAT_SECONDS) + yield f"data: {json.dumps(event, ensure_ascii=False)}\n\n" + except asyncio.TimeoutError: + yield ": ping\n\n" # heartbeat против idle-таймаутов прокси/туннеля + finally: + hub.unsubscribe(user.id, queue) # type: ignore[arg-type] + + return StreamingResponse( + stream(), + media_type="text/event-stream", + headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"}, + ) diff --git a/backend/app/routers/groups.py b/backend/app/routers/groups.py index 7813d16..7f61513 100644 --- a/backend/app/routers/groups.py +++ b/backend/app/routers/groups.py @@ -14,6 +14,7 @@ from app.services import ( group_service, invitation_service, membership_service, + notify, stats_service, ) @@ -153,6 +154,7 @@ def invite_member( user_agent=request.headers.get("user-agent"), ) session.commit() + notify.invitations_changed(invited.id) # type: ignore[arg-type] # живое появление у получателя return s.InvitationRead( id=inv.id, # type: ignore[arg-type] group_id=group_id, @@ -174,6 +176,7 @@ def remove_member( group_service.assert_member(session, group_id, user.id) # type: ignore[arg-type] group = group_service.get_group(session, group_id) membership_service.remove_member(session, group, user_id) + notify.group_changed(session, group_id, extra_user_ids=[user_id]) # + удалённому return s.OkResponse() diff --git a/backend/app/routers/invitations.py b/backend/app/routers/invitations.py index e388562..3f913c3 100644 --- a/backend/app/routers/invitations.py +++ b/backend/app/routers/invitations.py @@ -8,7 +8,7 @@ from app.auth.deps import get_current_user from app.db.session import get_session from app.models import User from app.schemas import api as s -from app.services import invitation_service +from app.services import invitation_service, notify router = APIRouter(prefix="/invitations", tags=["invitations"]) @@ -27,7 +27,9 @@ def accept_invitation( session: Session = Depends(get_session), user: User = Depends(get_current_user), ) -> s.OkResponse: - invitation_service.accept(session, user, invitation_id) + group = invitation_service.accept(session, user, invitation_id) + notify.invitations_changed(user.id) # type: ignore[arg-type] + notify.group_changed(session, group.id) # type: ignore[arg-type] # новый участник return s.OkResponse() @@ -38,4 +40,5 @@ def decline_invitation( user: User = Depends(get_current_user), ) -> s.OkResponse: invitation_service.decline(session, user, invitation_id) + notify.invitations_changed(user.id) # type: ignore[arg-type] return s.OkResponse() diff --git a/backend/app/routers/matches.py b/backend/app/routers/matches.py index ce79e28..7ee544d 100644 --- a/backend/app/routers/matches.py +++ b/backend/app/routers/matches.py @@ -1,7 +1,7 @@ """Роутер партий: рандом фракции, старт, завершение, детали, правка, удаление.""" from __future__ import annotations -from fastapi import APIRouter, Depends, File, Request, UploadFile +from fastapi import APIRouter, Depends, File, Query, Request, UploadFile from fastapi.responses import FileResponse from sqlmodel import Session @@ -16,6 +16,7 @@ from app.services import ( audit_service, group_service, match_service, + notify, user_service, ) from app.services.match_service import FinishInput, ParticipantInput, RosterInput @@ -69,6 +70,7 @@ def build_match_read(session: Session, match: Match, *, can_modify: bool = False overall_comment=match.overall_comment, created_by=match.created_by, can_modify=can_modify, + version=match_service.match_version(match), participants=parts, attachments=[ attachment_read(a, f"/api/matches/{match.id}") @@ -118,6 +120,7 @@ def start_match( user_agent=request.headers.get("user-agent"), ) session.commit() + notify.match_changed(session, match) return build_match_read(session, match, can_modify=match_service.can_modify(session, match, user)) @@ -143,6 +146,7 @@ def finish_match( win_reason=body.win_reason, overall_comment=body.overall_comment, overall_comment_set=("overall_comment" in body.model_fields_set), + expected_version=body.expected_version, ) audit_service.record( session, @@ -155,6 +159,7 @@ def finish_match( user_agent=request.headers.get("user-agent"), ) session.commit() + notify.match_changed(session, match) return build_match_read(session, match, can_modify=match_service.can_modify(session, match, user)) @@ -200,6 +205,7 @@ def update_match( win_reason=body.win_reason, win_reason_set=("win_reason" in body.model_fields_set), participants=participants, + expected_version=body.expected_version, ) audit_service.record( session, @@ -211,6 +217,7 @@ def update_match( user_agent=request.headers.get("user-agent"), ) session.commit() + notify.match_changed(session, match) return build_match_read(session, match, can_modify=match_service.can_modify(session, match, user)) @@ -240,6 +247,7 @@ def add_attachment( att = attachment_service.add_photo( session, match, user, content, ext, user_service.avatar_media_type(ext) ) + notify.match_changed(session, match) return attachment_read(att, f"/api/matches/{match_id}") @@ -253,6 +261,7 @@ def delete_attachment( match = match_service.get_match(session, match_id) _assert_can_attach(session, match, user) attachment_service.delete(session, match, attachment_id) + notify.match_changed(session, match) return s.OkResponse() @@ -278,6 +287,7 @@ def get_attachment( def delete_match( match_id: int, request: Request, + expected_version: str | None = Query(None), session: Session = Depends(get_session), user: User = Depends(get_current_user), ) -> s.OkResponse: @@ -285,7 +295,7 @@ def delete_match( match_service.assert_can_modify(session, match, user) match_id_val = match.id group_id_val = match.group_id - match_service.delete_match(session, match) + match_service.delete_match(session, match, expected_version=expected_version) audit_service.record( session, actor_id=user.id, @@ -297,4 +307,5 @@ def delete_match( user_agent=request.headers.get("user-agent"), ) session.commit() + notify.match_removed(session, match_id_val, group_id_val) # type: ignore[arg-type] return s.OkResponse() diff --git a/backend/app/schemas/api.py b/backend/app/schemas/api.py index cc6d9c0..c3b28db 100644 --- a/backend/app/schemas/api.py +++ b/backend/app/schemas/api.py @@ -192,6 +192,8 @@ class MatchFinish(BaseModel): participants: list[MatchFinishParticipant] win_reason: WinReason overall_comment: str | None = None + # Оптимистичная блокировка: версия партии, которую видел клиент (см. MatchRead.version). + expected_version: str | None = None # Полный участник (правка завершённой партии админом). @@ -208,6 +210,7 @@ class MatchUpdate(BaseModel): overall_comment: str | None = None win_reason: WinReason | None = None participants: list[ParticipantInput] | None = None + expected_version: str | None = None # оптимистичная блокировка class MatchParticipantRead(BaseModel): @@ -242,6 +245,7 @@ class MatchRead(BaseModel): overall_comment: str | None = None created_by: int can_modify: bool = False # может ли текущий зритель править/завершать партию + version: str # для оптимистичной блокировки (iso updated_at); клиент шлёт обратно participants: list[MatchParticipantRead] = [] attachments: list[AttachmentRead] = [] diff --git a/backend/app/services/match_service.py b/backend/app/services/match_service.py index 93abd29..f522670 100644 --- a/backend/app/services/match_service.py +++ b/backend/app/services/match_service.py @@ -16,7 +16,7 @@ from app.core.errors import ( NotFoundError, ValidationError, ) -from app.core.timeutil import app_today +from app.core.timeutil import app_today, iso_utc from app.models import Faction, GroupMember, Match, MatchParticipant, User from app.services import group_service @@ -70,6 +70,20 @@ def get_match(session: Session, match_id: int) -> Match: return match +def match_version(match: Match) -> str: + """Версия партии для оптимистичной блокировки (меняется при любом изменении).""" + return iso_utc(match.updated_at) + + +def assert_version(match: Match, expected: str | None) -> None: + """Если клиент прислал версию и она устарела — отказываем (кто-то изменил партию).""" + if expected is not None and expected != match_version(match): + raise ConflictError( + "Партия уже изменена на другом устройстве — обновите страницу.", + code="STALE_WRITE", + ) + + def participants_detail( session: Session, match_id: int ) -> list[tuple[MatchParticipant, User, Faction]]: @@ -193,7 +207,9 @@ def finish_match( win_reason: str, overall_comment: str | None = None, overall_comment_set: bool = False, + expected_version: str | None = None, ) -> Match: + assert_version(match, expected_version) if match.status != "in_progress": raise ConflictError("Партия уже завершена.") if win_reason not in WIN_REASONS: @@ -274,8 +290,10 @@ def update_match( win_reason: str | None = None, win_reason_set: bool = False, participants: list[ParticipantInput] | None = None, + expected_version: str | None = None, ) -> Match: """Правка завершённой партии (админ): полный список участников с местами.""" + assert_version(match, expected_version) if played_at is not None: match.played_at = played_at if overall_comment_set: @@ -317,9 +335,10 @@ def update_match( return match -def delete_match(session: Session, match: Match) -> None: +def delete_match(session: Session, match: Match, expected_version: str | None = None) -> None: from app.services import attachment_service # избегаем цикла импорта + assert_version(match, expected_version) match_id = match.id session.delete(match) # участники и вложения (БД) удалятся каскадом (FK ON DELETE CASCADE) session.commit() diff --git a/backend/app/services/notify.py b/backend/app/services/notify.py new file mode 100644 index 0000000..b39cc04 --- /dev/null +++ b/backend/app/services/notify.py @@ -0,0 +1,46 @@ +"""Публикация SSE-событий «данные изменились» нужным пользователям. + +Тонкая прослойка над core.events.hub — держит роутеры чистыми. Событие несёт лишь тип и id, +клиент по нему точечно перезапрашивает данные. +""" +from __future__ import annotations + +from sqlmodel import Session, select + +from app.core.events import hub +from app.models import GroupMember, Match + + +def _group_member_ids(session: Session, group_id: int) -> list[int]: + return list( + session.exec(select(GroupMember.user_id).where(GroupMember.group_id == group_id)).all() + ) + + +def match_changed(session: Session, match: Match) -> None: + """Партия изменилась — уведомить всех участников её группы.""" + hub.publish( + _group_member_ids(session, match.group_id), + {"type": "match", "match_id": match.id, "group_id": match.group_id}, + ) + + +def match_removed(session: Session, match_id: int, group_id: int) -> None: + """Партия удалена — уведомить участников группы (обновить списки).""" + hub.publish( + _group_member_ids(session, group_id), + {"type": "match", "match_id": match_id, "group_id": group_id}, + ) + + +def group_changed(session: Session, group_id: int, extra_user_ids: list[int] | None = None) -> None: + """Состав/данные группы изменились — участникам (+ доп. адресатам, напр. удалённому).""" + ids = _group_member_ids(session, group_id) + if extra_user_ids: + ids = ids + extra_user_ids + hub.publish(ids, {"type": "group", "group_id": group_id}) + + +def invitations_changed(user_id: int) -> None: + """У пользователя изменился список приглашений.""" + hub.publish([user_id], {"type": "invitations"}) diff --git a/backend/tests/test_concurrency.py b/backend/tests/test_concurrency.py new file mode 100644 index 0000000..516621f --- /dev/null +++ b/backend/tests/test_concurrency.py @@ -0,0 +1,77 @@ +"""Оптимистичная блокировка партии: устаревшие правки/удаление отклоняются (STALE_WRITE).""" +from __future__ import annotations + +from fastapi.testclient import TestClient + +from tests.conftest import add_group_member, csrf_headers, finish_match, login, start_match + + +def _start(client: TestClient, engine) -> tuple[dict, int, int]: + me = login(client, "Хост") + exps = [e["id"] for e in client.get("/api/expansions").json()] + gid = client.post( + "/api/groups", json={"name": "Группа", "expansion_ids": exps}, headers=csrf_headers(client) + ).json()["id"] + p2 = add_group_member(engine, gid, "Игрок2") + fids = [f["id"] for f in client.get(f"/api/groups/{gid}/factions").json()] + started = start_match( + client, gid, + [{"user_id": me["id"], "faction_id": fids[0]}, {"user_id": p2, "faction_id": fids[1]}], + ) + assert started.status_code == 200, started.text + return me, p2, started.json()["id"] + + +def test_stale_delete_rejected(client: TestClient, engine): + """Сценарий бага: ПК завершил, телефон со старой версией жмёт «Отменить».""" + me, p2, mid = _start(client, engine) + v1 = client.get(f"/api/matches/{mid}").json()["version"] + + # «ПК» завершает партию — версия меняется. + fin = finish_match( + client, mid, [{"user_id": me["id"], "place": 1}, {"user_id": p2, "place": 2}], + win_reason="objectives", + ) + assert fin.status_code == 200, fin.text + + # «Телефон» со старой версией пытается отменить → 409 STALE_WRITE, партия НЕ удаляется. + stale = client.delete( + f"/api/matches/{mid}", params={"expected_version": v1}, headers=csrf_headers(client) + ) + assert stale.status_code == 409, stale.text + assert stale.json()["error"]["code"] == "STALE_WRITE" + assert client.get(f"/api/matches/{mid}").status_code == 200 # жива + + # С актуальной версией удаление проходит. + v2 = client.get(f"/api/matches/{mid}").json()["version"] + ok = client.delete( + f"/api/matches/{mid}", params={"expected_version": v2}, headers=csrf_headers(client) + ) + assert ok.status_code == 200, ok.text + assert client.get(f"/api/matches/{mid}").status_code == 404 + + +def test_stale_finish_rejected(client: TestClient, engine): + me, p2, mid = _start(client, engine) + v1 = client.get(f"/api/matches/{mid}").json()["version"] + + # Партию изменили (правка комментария) — версия устарела. + bump = client.patch( + f"/api/matches/{mid}", json={"overall_comment": "правка"}, headers=csrf_headers(client) + ) + assert bump.status_code == 200, bump.text + + # Завершение со старой версией → 409 STALE_WRITE. + r = client.post( + f"/api/matches/{mid}/finish", + json={ + "participants": [ + {"user_id": me["id"], "place": 1}, + {"user_id": p2, "place": 2}, + ], + "win_reason": "objectives", + "expected_version": v1, + }, + headers=csrf_headers(client), + ) + assert r.status_code == 409 and r.json()["error"]["code"] == "STALE_WRITE", r.text diff --git a/backend/tests/test_core_flow.py b/backend/tests/test_core_flow.py index 4a6f407..e55ea9e 100644 --- a/backend/tests/test_core_flow.py +++ b/backend/tests/test_core_flow.py @@ -287,6 +287,19 @@ def test_disabled_account_cannot_login(client: TestClient, make_admin): assert guest_dev["is_active"] is False +def test_csrf_required_for_session_mutations(client: TestClient): + """После рефактора CSRF на ASGI защита сохраняется: мутация с сессией без X-CSRF-Token → 403.""" + login(client, "Аня") # появились cookie сессии и csrf_token + no_header = client.post("/api/groups", json={"name": "Группа", "expansion_ids": []}) + assert no_header.status_code == 403 + assert no_header.json()["error"]["code"] == "CSRF_FAILED" + # С корректным заголовком — проходит. + ok = client.post( + "/api/groups", json={"name": "Группа", "expansion_ids": []}, headers=csrf_headers(client) + ) + assert ok.status_code == 200, ok.text + + def test_group_stats_includes_inactive_members(client: TestClient, engine): """Участники без завершённых партий попадают в отдельный блок inactive (не в provisional).""" me = login(client, "Капитан") diff --git a/backend/tests/test_events.py b/backend/tests/test_events.py new file mode 100644 index 0000000..e43f351 --- /dev/null +++ b/backend/tests/test_events.py @@ -0,0 +1,43 @@ +"""SSE-шина: доставка событий подписчику; эндпойнт /api/events требует вход.""" +from __future__ import annotations + +import asyncio + +from fastapi.testclient import TestClient + + +def test_hub_delivers_to_subscriber(): + from app.core.events import EventHub + + async def run(): + hub = EventHub() + hub.bind_loop(asyncio.get_running_loop()) + q = hub.subscribe(1) + hub.publish([1, 2], {"type": "match", "match_id": 5}) + event = await asyncio.wait_for(q.get(), timeout=1) + assert event == {"type": "match", "match_id": 5} + hub.unsubscribe(1, q) + + asyncio.run(run()) + + +def test_hub_isolates_users(): + from app.core.events import EventHub + + async def run(): + hub = EventHub() + hub.bind_loop(asyncio.get_running_loop()) + q1 = hub.subscribe(1) + q2 = hub.subscribe(2) + hub.publish([2], {"type": "invitations"}) # только пользователю 2 + got2 = await asyncio.wait_for(q2.get(), timeout=1) + assert got2 == {"type": "invitations"} + assert q1.empty() # пользователю 1 ничего не пришло + + asyncio.run(run()) + + +def test_events_requires_auth(client: TestClient): + # Без сессии SSE-эндпойнт не отдаёт поток (401 на зависимости get_current_user). + r = client.get("/api/events") + assert r.status_code == 401 diff --git a/deploy/README.md b/deploy/README.md index e39409e..fe335ca 100644 --- a/deploy/README.md +++ b/deploy/README.md @@ -45,3 +45,15 @@ HTTPS твоими сертификатами и проксирует трафи > Telegram-вход требует HTTPS-домен: у BotFather `/setdomain` укажи оба домена > (`forbiddenstars.ru` и `forbidden-stars.ru`). + +## Реал-тайм (SSE) + +Приложение шлёт события «данные изменились» через SSE: `GET /api/events` +(`text/event-stream`). Это обычный HTTP — проходит через Caddy и SSH-туннель без доп. +настроек, КРОМЕ одного: для пути `/api/events` в `Caddyfile` выделен отдельный `handle` +**без `encode`** и с **`flush_interval -1`** (иначе сжатие/буферизация задержат поток). +После правки `Caddyfile` на VPS: `sudo systemctl reload caddy` (или `restart`). + +Шина событий — **in-memory**, рассчитана на один процесс (uvicorn `--workers 1`, как в +контейнере). Если когда-нибудь поднимешь несколько воркеров/реплик — шину нужно вынести во +внешний брокер (Redis pub/sub), иначе события увидит только тот воркер, что принял мутацию. diff --git a/deploy/vps/Caddyfile b/deploy/vps/Caddyfile index cc66ff5..e7f360b 100644 --- a/deploy/vps/Caddyfile +++ b/deploy/vps/Caddyfile @@ -8,15 +8,34 @@ # приложению добавляет X-Forwarded-Proto=https / X-Forwarded-For / Host — # приложение это учитывает (uvicorn --proxy-headers). Положи файл в /etc/caddy/Caddyfile. # Сертификаты — см. deploy/vps/README.md (fullchain = leaf + промежуточные одним файлом). +# +# SSE (/api/events): отдельный handle БЕЗ encode и с flush_interval -1 — иначе сжатие/ +# буферизация задержат поток событий (реал-тайм перестанет работать). forbiddenstars.ru { tls /etc/caddy/certs/forbiddenstars.ru/fullchain.pem /etc/caddy/certs/forbiddenstars.ru/privkey.pem - encode zstd gzip - reverse_proxy 127.0.0.1:9000 + @sse path /api/events + handle @sse { + reverse_proxy 127.0.0.1:9000 { + flush_interval -1 + } + } + handle { + encode zstd gzip + reverse_proxy 127.0.0.1:9000 + } } forbidden-stars.ru { tls /etc/caddy/certs/forbidden-stars.ru/fullchain.pem /etc/caddy/certs/forbidden-stars.ru/privkey.pem - encode zstd gzip - reverse_proxy 127.0.0.1:9001 + @sse path /api/events + handle @sse { + reverse_proxy 127.0.0.1:9001 { + flush_interval -1 + } + } + handle { + encode zstd gzip + reverse_proxy 127.0.0.1:9001 + } } diff --git a/frontend/src/api/schema.d.ts b/frontend/src/api/schema.d.ts index a309b3c..0535004 100644 --- a/frontend/src/api/schema.d.ts +++ b/frontend/src/api/schema.d.ts @@ -645,6 +645,23 @@ export interface paths { patch?: never; trace?: never; }; + "/api/events": { + parameters: { + query?: never; + header?: never; + path?: never; + cookie?: never; + }; + /** Events */ + get: operations["events_api_events_get"]; + put?: never; + post?: never; + delete?: never; + options?: never; + head?: never; + patch?: never; + trace?: never; + }; "/api/admin/auth/login": { parameters: { query?: never; @@ -1507,6 +1524,8 @@ export interface components { win_reason: "objectives" | "worlds" | "plastic" | "resources"; /** Overall Comment */ overall_comment?: string | null; + /** Expected Version */ + expected_version?: string | null; }; /** MatchFinishParticipant */ MatchFinishParticipant: { @@ -1627,6 +1646,8 @@ export interface components { * @default false */ can_modify: boolean; + /** Version */ + version: string; /** * Participants * @default [] @@ -1648,6 +1669,8 @@ export interface components { win_reason?: ("objectives" | "worlds" | "plastic" | "resources") | null; /** Participants */ participants?: components["schemas"]["ParticipantInput"][] | null; + /** Expected Version */ + expected_version?: string | null; }; /** MeRead */ MeRead: { @@ -2882,7 +2905,9 @@ export interface operations { }; delete_match_api_matches__match_id__delete: { parameters: { - query?: never; + query?: { + expected_version?: string | null; + }; header?: never; path: { match_id: number; @@ -3187,6 +3212,26 @@ export interface operations { }; }; }; + events_api_events_get: { + parameters: { + query?: never; + header?: never; + path?: never; + cookie?: never; + }; + requestBody?: never; + responses: { + /** @description Successful Response */ + 200: { + headers: { + [name: string]: unknown; + }; + content: { + "application/json": unknown; + }; + }; + }; + }; admin_login_api_admin_auth_login_post: { parameters: { query?: never; diff --git a/frontend/src/app/queryClient.ts b/frontend/src/app/queryClient.ts index e6ec224..727d29b 100644 --- a/frontend/src/app/queryClient.ts +++ b/frontend/src/app/queryClient.ts @@ -4,8 +4,11 @@ export const queryClient = new QueryClient({ defaultOptions: { queries: { retry: false, - refetchOnWindowFocus: false, - staleTime: 30_000, + // Свежесть: возврат во вкладку/восстановление сети → перезапрос (на случай, если + // SSE-соединение временно отвалилось). Точечные апдейты приходят пушем (useServerEvents). + refetchOnWindowFocus: true, + refetchOnReconnect: true, + staleTime: 15_000, }, }, }); diff --git a/frontend/src/components/AppShell.tsx b/frontend/src/components/AppShell.tsx index 1fed4ed..dc40575 100644 --- a/frontend/src/components/AppShell.tsx +++ b/frontend/src/components/AppShell.tsx @@ -1,5 +1,7 @@ import { Outlet, useLocation, useNavigate } from "react-router-dom"; +import { useMe } from "../hooks/auth"; +import { useServerEvents } from "../hooks/useServerEvents"; import { BottomBar } from "./BottomBar"; const TITLES: Record = { @@ -14,6 +16,8 @@ const TITLES: Record = { export function AppShell() { const location = useLocation(); const navigate = useNavigate(); + const { data: me } = useMe(); + useServerEvents(!!me); // живые обновления (SSE) только при наличии сессии const title = TITLES[location.pathname] ?? (location.pathname.startsWith("/u/") ? "Профиль игрока" : "Forbidden Stars"); diff --git a/frontend/src/hooks/matches.ts b/frontend/src/hooks/matches.ts index 75f9949..90e2b94 100644 --- a/frontend/src/hooks/matches.ts +++ b/frontend/src/hooks/matches.ts @@ -70,10 +70,14 @@ export function useFinishMatch() { export function useDeleteMatch() { const qc = useQueryClient(); return useMutation({ - mutationFn: async (matchId: number) => + // expectedVersion — оптимистичная блокировка: отмена устаревшей версии вернёт 409 STALE_WRITE. + mutationFn: async (args: { matchId: number; expectedVersion?: string }) => unwrap( await api.DELETE("/api/matches/{match_id}", { - params: { path: { match_id: matchId } }, + params: { + path: { match_id: args.matchId }, + query: args.expectedVersion ? { expected_version: args.expectedVersion } : {}, + }, }), ), onSuccess: () => qc.invalidateQueries(), diff --git a/frontend/src/hooks/useServerEvents.ts b/frontend/src/hooks/useServerEvents.ts new file mode 100644 index 0000000..8b1f554 --- /dev/null +++ b/frontend/src/hooks/useServerEvents.ts @@ -0,0 +1,56 @@ +import { useQueryClient } from "@tanstack/react-query"; +import { useEffect } from "react"; + +import { qk } from "../api/queryKeys"; + +interface ServerEvent { + type: "match" | "group" | "invitations"; + match_id?: number; + group_id?: number; +} + +/** + * Подписка на SSE-поток `/api/events` (push-to-invalidate): по событию с сервера точечно + * инвалидируем нужные запросы, и TanStack Query сам их перезапрашивает. EventSource + * переподключается автоматически. Открываем одно соединение, только когда авторизованы. + */ +export function useServerEvents(enabled: boolean) { + const qc = useQueryClient(); + useEffect(() => { + if (!enabled) return; + const base = import.meta.env.VITE_API_BASE_URL || ""; + const es = new EventSource(`${base}/api/events`, { withCredentials: true }); + + es.onmessage = (e) => { + let ev: ServerEvent; + try { + ev = JSON.parse(e.data) as ServerEvent; + } catch { + return; + } + if (ev.type === "invitations") { + qc.invalidateQueries({ queryKey: qk.invitations }); + } else if (ev.type === "match") { + if (ev.match_id != null) qc.invalidateQueries({ queryKey: qk.match(ev.match_id) }); + if (ev.group_id != null) { + qc.invalidateQueries({ queryKey: qk.groupMatches(ev.group_id) }); + qc.invalidateQueries({ queryKey: qk.groupStats(ev.group_id) }); + } + qc.invalidateQueries({ queryKey: qk.home }); + qc.invalidateQueries({ queryKey: qk.leaderboard }); + } else if (ev.type === "group") { + if (ev.group_id != null) { + qc.invalidateQueries({ queryKey: qk.group(ev.group_id) }); + qc.invalidateQueries({ queryKey: qk.groupMembers(ev.group_id) }); + qc.invalidateQueries({ queryKey: qk.groupStats(ev.group_id) }); + } + qc.invalidateQueries({ queryKey: qk.me }); + qc.invalidateQueries({ queryKey: qk.home }); + qc.invalidateQueries({ queryKey: qk.invitations }); + } + }; + // onerror не логируем: EventSource переподключается сам. + + return () => es.close(); + }, [enabled, qc]); +} diff --git a/frontend/src/pages/MatchDetailPage.tsx b/frontend/src/pages/MatchDetailPage.tsx index 3232152..e72d97d 100644 --- a/frontend/src/pages/MatchDetailPage.tsx +++ b/frontend/src/pages/MatchDetailPage.tsx @@ -26,7 +26,7 @@ interface FinishRow { export function MatchDetailPage() { const { matchId } = useParams(); const id = matchId ? Number(matchId) : null; - const { data: match, isLoading } = useMatch(id); + const { data: match, isLoading, refetch } = useMatch(id); const finish = useFinishMatch(); const del = useDeleteMatch(); const uploadAtt = useUploadMatchAttachment(id ?? 0); @@ -59,8 +59,11 @@ export function MatchDetailPage() { const upd = (i: number, patch: Partial) => setRows(finishRows.map((r, idx) => (idx === i ? { ...r, ...patch } : r))); + // Конфликт версий (кто-то изменил партию с другого устройства) → сообщаем и обновляем. + const isStale = (e: unknown) => e instanceof ApiError && e.code === "STALE_WRITE"; + const submitFinish = async () => { - if (!id) return; + if (!id || !match) return; setError(null); try { await finish.mutateAsync({ @@ -73,19 +76,34 @@ export function MatchDetailPage() { })), win_reason: winReason, overall_comment: overall.trim() || null, + expected_version: match.version, }, }); toast.show("Партия завершена"); } catch (e) { + if (isStale(e)) { + toast.show("Партия изменилась на другом устройстве — обновлено"); + refetch(); + return; + } setError(e instanceof ApiError ? e.message : "Не удалось завершить партию"); } }; const remove = async () => { - if (!id) return; - await del.mutateAsync(id).catch(() => {}); - toast.show(inProgress ? "Партия отменена" : "Партия удалена"); - navigate(-1); + if (!id || !match) return; + try { + await del.mutateAsync({ matchId: id, expectedVersion: match.version }); + toast.show(inProgress ? "Партия отменена" : "Партия удалена"); + navigate(-1); + } catch (e) { + if (isStale(e)) { + toast.show("Партия изменилась на другом устройстве — обновлено"); + refetch(); + return; + } + toast.show(e instanceof ApiError ? e.message : "Не удалось выполнить"); + } }; const sorted = [...match.participants].sort(