v0.6 - Настройка SSE-соединения для пуша данных с сервера

This commit is contained in:
2026-06-17 23:22:09 +03:00
21 changed files with 553 additions and 48 deletions
+56
View File
@@ -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()
+44 -27
View File
@@ -10,7 +10,6 @@ from fastapi import FastAPI, Request
from fastapi.exceptions import RequestValidationError from fastapi.exceptions import RequestValidationError
from fastapi.middleware.cors import CORSMiddleware from fastapi.middleware.cors import CORSMiddleware
from fastapi.responses import FileResponse, JSONResponse from fastapi.responses import FileResponse, JSONResponse
from starlette.middleware.base import BaseHTTPMiddleware
from app.core import security from app.core import security
from app.core.config import settings from app.core.config import settings
@@ -19,6 +18,7 @@ from app.routers import (
achievements, achievements,
admin, admin,
auth, auth,
events,
groups, groups,
invitations, invitations,
matches, matches,
@@ -33,36 +33,53 @@ _STATIC_DIR = Path(os.getenv("STATIC_DIR", str(Path(__file__).resolve().parent.p
_UNSAFE_METHODS = {"POST", "PUT", "PATCH", "DELETE"} _UNSAFE_METHODS = {"POST", "PUT", "PATCH", "DELETE"}
class CSRFMiddleware(BaseHTTPMiddleware): class CSRFMiddleware:
"""Double-submit CSRF: для аутентифицированных мутаций на /api требуем """Double-submit CSRF на чистом ASGI: для аутентифицированных мутаций на /api требуем
совпадения заголовка X-CSRF-Token и cookie csrf_token.""" совпадения заголовка X-CSRF-Token и cookie csrf_token.
async def dispatch(self, request: Request, call_next): # noqa: ANN001 Намеренно НЕ на BaseHTTPMiddleware: тот буферизует потоковые ответы и ломает SSE
path = request.url.path (/api/events). Чистый ASGI пропускает стримы насквозь, вмешиваясь только при отказе CSRF.
if request.method in _UNSAFE_METHODS and path.startswith("/api"): """
has_session = (
security.USER_COOKIE in request.cookies def __init__(self, app) -> None: # noqa: ANN001
or security.ADMIN_COOKIE in request.cookies self.app = app
)
if has_session: async def __call__(self, scope, receive, send): # noqa: ANN001
cookie_token = request.cookies.get(security.CSRF_COOKIE) if scope["type"] == "http":
header_token = request.headers.get(security.CSRF_HEADER) request = Request(scope)
if not cookie_token or cookie_token != header_token: if request.method in _UNSAFE_METHODS and request.url.path.startswith("/api"):
return JSONResponse( has_session = (
status_code=403, security.USER_COOKIE in request.cookies
content={ or security.ADMIN_COOKIE in request.cookies
"error": { )
"code": "CSRF_FAILED", if has_session:
"message": "Неверный или отсутствующий CSRF-токен.", cookie_token = request.cookies.get(security.CSRF_COOKIE)
"details": None, header_token = request.headers.get(security.CSRF_HEADER)
} if not cookie_token or cookie_token != header_token:
}, response = JSONResponse(
) status_code=403,
return await call_next(request) content={
"error": {
"code": "CSRF_FAILED",
"message": "Неверный или отсутствующий CSRF-токен.",
"details": None,
}
},
)
await response(scope, receive, send)
return
await self.app(scope, receive, send)
@asynccontextmanager @asynccontextmanager
async def _lifespan(_app: FastAPI): 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 при старте # В DEV приложение само подтягивает справочники и админа из .env при старте
# (в test/prod это делает entrypoint.sh; в pytest отключено FS_STARTUP_BOOTSTRAP=0). # (в test/prod это делает entrypoint.sh; в pytest отключено FS_STARTUP_BOOTSTRAP=0).
if settings.is_development and os.getenv("FS_STARTUP_BOOTSTRAP", "1") != "0": if settings.is_development and os.getenv("FS_STARTUP_BOOTSTRAP", "1") != "0":
@@ -118,7 +135,7 @@ def create_app() -> FastAPI:
# API-роутеры под /api. # API-роутеры под /api.
api_routers = [auth.router, users.router, groups.router, invitations.router, api_routers = [auth.router, users.router, groups.router, invitations.router,
matches.router, reference.router, stats.router, achievements.router, matches.router, reference.router, stats.router, achievements.router,
admin.router] events.router, admin.router]
for r in api_routers: for r in api_routers:
app.include_router(r, prefix="/api") app.include_router(r, prefix="/api")
+7
View File
@@ -20,6 +20,7 @@ from app.services import (
audit_service, audit_service,
faction_service, faction_service,
match_service, match_service,
notify,
user_service, user_service,
) )
@@ -229,6 +230,7 @@ def update_match(
win_reason=body.win_reason, win_reason=body.win_reason,
win_reason_set=("win_reason" in body.model_fields_set), win_reason_set=("win_reason" in body.model_fields_set),
participants=participants, participants=participants,
expected_version=body.expected_version,
) )
audit_service.record( audit_service.record(
session, session,
@@ -239,6 +241,7 @@ def update_match(
ip=request.client.host if request.client else None, ip=request.client.host if request.client else None,
) )
session.commit() session.commit()
notify.match_changed(session, match)
return build_match_read(session, match, can_modify=True) return build_match_read(session, match, can_modify=True)
@@ -288,12 +291,14 @@ def delete_match(
session: Session = Depends(get_session), session: Session = Depends(get_session),
admin: User = Depends(get_current_admin), admin: User = Depends(get_current_admin),
) -> s.OkResponse: ) -> s.OkResponse:
group_id = match_service.get_match(session, match_id).group_id # для уведомления
admin_service.delete_match(session, match_id) admin_service.delete_match(session, match_id)
audit_service.record( audit_service.record(
session, actor_id=admin.id, action="delete", entity_type="match", entity_id=match_id, session, actor_id=admin.id, action="delete", entity_type="match", entity_id=match_id,
ip=request.client.host if request.client else None, ip=request.client.host if request.client else None,
) )
session.commit() session.commit()
notify.match_removed(session, match_id, group_id)
return s.OkResponse() return s.OkResponse()
@@ -329,6 +334,7 @@ def admin_add_attachment(
att = attachment_service.add_photo( att = attachment_service.add_photo(
session, match, admin, content, ext, user_service.avatar_media_type(ext) 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}") return attachment_read(att, f"/api/admin/matches/{match_id}")
@@ -341,6 +347,7 @@ def admin_delete_attachment(
) -> s.OkResponse: ) -> s.OkResponse:
match = match_service.get_match(session, match_id) match = match_service.get_match(session, match_id)
attachment_service.delete(session, match, attachment_id) attachment_service.delete(session, match, attachment_id)
notify.match_changed(session, match)
return s.OkResponse() return s.OkResponse()
+45
View File
@@ -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"},
)
+3
View File
@@ -14,6 +14,7 @@ from app.services import (
group_service, group_service,
invitation_service, invitation_service,
membership_service, membership_service,
notify,
stats_service, stats_service,
) )
@@ -153,6 +154,7 @@ def invite_member(
user_agent=request.headers.get("user-agent"), user_agent=request.headers.get("user-agent"),
) )
session.commit() session.commit()
notify.invitations_changed(invited.id) # type: ignore[arg-type] # живое появление у получателя
return s.InvitationRead( return s.InvitationRead(
id=inv.id, # type: ignore[arg-type] id=inv.id, # type: ignore[arg-type]
group_id=group_id, 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_service.assert_member(session, group_id, user.id) # type: ignore[arg-type]
group = group_service.get_group(session, group_id) group = group_service.get_group(session, group_id)
membership_service.remove_member(session, group, user_id) membership_service.remove_member(session, group, user_id)
notify.group_changed(session, group_id, extra_user_ids=[user_id]) # + удалённому
return s.OkResponse() return s.OkResponse()
+5 -2
View File
@@ -8,7 +8,7 @@ from app.auth.deps import get_current_user
from app.db.session import get_session from app.db.session import get_session
from app.models import User from app.models import User
from app.schemas import api as s 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"]) router = APIRouter(prefix="/invitations", tags=["invitations"])
@@ -27,7 +27,9 @@ def accept_invitation(
session: Session = Depends(get_session), session: Session = Depends(get_session),
user: User = Depends(get_current_user), user: User = Depends(get_current_user),
) -> s.OkResponse: ) -> 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() return s.OkResponse()
@@ -38,4 +40,5 @@ def decline_invitation(
user: User = Depends(get_current_user), user: User = Depends(get_current_user),
) -> s.OkResponse: ) -> s.OkResponse:
invitation_service.decline(session, user, invitation_id) invitation_service.decline(session, user, invitation_id)
notify.invitations_changed(user.id) # type: ignore[arg-type]
return s.OkResponse() return s.OkResponse()
+13 -2
View File
@@ -1,7 +1,7 @@
"""Роутер партий: рандом фракции, старт, завершение, детали, правка, удаление.""" """Роутер партий: рандом фракции, старт, завершение, детали, правка, удаление."""
from __future__ import annotations 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 fastapi.responses import FileResponse
from sqlmodel import Session from sqlmodel import Session
@@ -16,6 +16,7 @@ from app.services import (
audit_service, audit_service,
group_service, group_service,
match_service, match_service,
notify,
user_service, user_service,
) )
from app.services.match_service import FinishInput, ParticipantInput, RosterInput 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, overall_comment=match.overall_comment,
created_by=match.created_by, created_by=match.created_by,
can_modify=can_modify, can_modify=can_modify,
version=match_service.match_version(match),
participants=parts, participants=parts,
attachments=[ attachments=[
attachment_read(a, f"/api/matches/{match.id}") attachment_read(a, f"/api/matches/{match.id}")
@@ -118,6 +120,7 @@ def start_match(
user_agent=request.headers.get("user-agent"), user_agent=request.headers.get("user-agent"),
) )
session.commit() session.commit()
notify.match_changed(session, match)
return build_match_read(session, match, can_modify=match_service.can_modify(session, match, user)) 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, win_reason=body.win_reason,
overall_comment=body.overall_comment, overall_comment=body.overall_comment,
overall_comment_set=("overall_comment" in body.model_fields_set), overall_comment_set=("overall_comment" in body.model_fields_set),
expected_version=body.expected_version,
) )
audit_service.record( audit_service.record(
session, session,
@@ -155,6 +159,7 @@ def finish_match(
user_agent=request.headers.get("user-agent"), user_agent=request.headers.get("user-agent"),
) )
session.commit() session.commit()
notify.match_changed(session, match)
return build_match_read(session, match, can_modify=match_service.can_modify(session, match, user)) 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=body.win_reason,
win_reason_set=("win_reason" in body.model_fields_set), win_reason_set=("win_reason" in body.model_fields_set),
participants=participants, participants=participants,
expected_version=body.expected_version,
) )
audit_service.record( audit_service.record(
session, session,
@@ -211,6 +217,7 @@ def update_match(
user_agent=request.headers.get("user-agent"), user_agent=request.headers.get("user-agent"),
) )
session.commit() session.commit()
notify.match_changed(session, match)
return build_match_read(session, match, can_modify=match_service.can_modify(session, match, user)) 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( att = attachment_service.add_photo(
session, match, user, content, ext, user_service.avatar_media_type(ext) 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}") return attachment_read(att, f"/api/matches/{match_id}")
@@ -253,6 +261,7 @@ def delete_attachment(
match = match_service.get_match(session, match_id) match = match_service.get_match(session, match_id)
_assert_can_attach(session, match, user) _assert_can_attach(session, match, user)
attachment_service.delete(session, match, attachment_id) attachment_service.delete(session, match, attachment_id)
notify.match_changed(session, match)
return s.OkResponse() return s.OkResponse()
@@ -278,6 +287,7 @@ def get_attachment(
def delete_match( def delete_match(
match_id: int, match_id: int,
request: Request, request: Request,
expected_version: str | None = Query(None),
session: Session = Depends(get_session), session: Session = Depends(get_session),
user: User = Depends(get_current_user), user: User = Depends(get_current_user),
) -> s.OkResponse: ) -> s.OkResponse:
@@ -285,7 +295,7 @@ def delete_match(
match_service.assert_can_modify(session, match, user) match_service.assert_can_modify(session, match, user)
match_id_val = match.id match_id_val = match.id
group_id_val = match.group_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( audit_service.record(
session, session,
actor_id=user.id, actor_id=user.id,
@@ -297,4 +307,5 @@ def delete_match(
user_agent=request.headers.get("user-agent"), user_agent=request.headers.get("user-agent"),
) )
session.commit() session.commit()
notify.match_removed(session, match_id_val, group_id_val) # type: ignore[arg-type]
return s.OkResponse() return s.OkResponse()
+4
View File
@@ -192,6 +192,8 @@ class MatchFinish(BaseModel):
participants: list[MatchFinishParticipant] participants: list[MatchFinishParticipant]
win_reason: WinReason win_reason: WinReason
overall_comment: str | None = None overall_comment: str | None = None
# Оптимистичная блокировка: версия партии, которую видел клиент (см. MatchRead.version).
expected_version: str | None = None
# Полный участник (правка завершённой партии админом). # Полный участник (правка завершённой партии админом).
@@ -208,6 +210,7 @@ class MatchUpdate(BaseModel):
overall_comment: str | None = None overall_comment: str | None = None
win_reason: WinReason | None = None win_reason: WinReason | None = None
participants: list[ParticipantInput] | None = None participants: list[ParticipantInput] | None = None
expected_version: str | None = None # оптимистичная блокировка
class MatchParticipantRead(BaseModel): class MatchParticipantRead(BaseModel):
@@ -242,6 +245,7 @@ class MatchRead(BaseModel):
overall_comment: str | None = None overall_comment: str | None = None
created_by: int created_by: int
can_modify: bool = False # может ли текущий зритель править/завершать партию can_modify: bool = False # может ли текущий зритель править/завершать партию
version: str # для оптимистичной блокировки (iso updated_at); клиент шлёт обратно
participants: list[MatchParticipantRead] = [] participants: list[MatchParticipantRead] = []
attachments: list[AttachmentRead] = [] attachments: list[AttachmentRead] = []
+21 -2
View File
@@ -16,7 +16,7 @@ from app.core.errors import (
NotFoundError, NotFoundError,
ValidationError, 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.models import Faction, GroupMember, Match, MatchParticipant, User
from app.services import group_service from app.services import group_service
@@ -70,6 +70,20 @@ def get_match(session: Session, match_id: int) -> Match:
return 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( def participants_detail(
session: Session, match_id: int session: Session, match_id: int
) -> list[tuple[MatchParticipant, User, Faction]]: ) -> list[tuple[MatchParticipant, User, Faction]]:
@@ -193,7 +207,9 @@ def finish_match(
win_reason: str, win_reason: str,
overall_comment: str | None = None, overall_comment: str | None = None,
overall_comment_set: bool = False, overall_comment_set: bool = False,
expected_version: str | None = None,
) -> Match: ) -> Match:
assert_version(match, expected_version)
if match.status != "in_progress": if match.status != "in_progress":
raise ConflictError("Партия уже завершена.") raise ConflictError("Партия уже завершена.")
if win_reason not in WIN_REASONS: if win_reason not in WIN_REASONS:
@@ -274,8 +290,10 @@ def update_match(
win_reason: str | None = None, win_reason: str | None = None,
win_reason_set: bool = False, win_reason_set: bool = False,
participants: list[ParticipantInput] | None = None, participants: list[ParticipantInput] | None = None,
expected_version: str | None = None,
) -> Match: ) -> Match:
"""Правка завершённой партии (админ): полный список участников с местами.""" """Правка завершённой партии (админ): полный список участников с местами."""
assert_version(match, expected_version)
if played_at is not None: if played_at is not None:
match.played_at = played_at match.played_at = played_at
if overall_comment_set: if overall_comment_set:
@@ -317,9 +335,10 @@ def update_match(
return 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 # избегаем цикла импорта from app.services import attachment_service # избегаем цикла импорта
assert_version(match, expected_version)
match_id = match.id match_id = match.id
session.delete(match) # участники и вложения (БД) удалятся каскадом (FK ON DELETE CASCADE) session.delete(match) # участники и вложения (БД) удалятся каскадом (FK ON DELETE CASCADE)
session.commit() session.commit()
+46
View File
@@ -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"})
+77
View File
@@ -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
+13
View File
@@ -287,6 +287,19 @@ def test_disabled_account_cannot_login(client: TestClient, make_admin):
assert guest_dev["is_active"] is False 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): def test_group_stats_includes_inactive_members(client: TestClient, engine):
"""Участники без завершённых партий попадают в отдельный блок inactive (не в provisional).""" """Участники без завершённых партий попадают в отдельный блок inactive (не в provisional)."""
me = login(client, "Капитан") me = login(client, "Капитан")
+43
View File
@@ -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
+12
View File
@@ -45,3 +45,15 @@ HTTPS твоими сертификатами и проксирует трафи
> Telegram-вход требует HTTPS-домен: у BotFather `/setdomain` укажи оба домена > Telegram-вход требует HTTPS-домен: у BotFather `/setdomain` укажи оба домена
> (`forbiddenstars.ru` и `forbidden-stars.ru`). > (`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), иначе события увидит только тот воркер, что принял мутацию.
+23 -4
View File
@@ -8,15 +8,34 @@
# приложению добавляет X-Forwarded-Proto=https / X-Forwarded-For / Host — # приложению добавляет X-Forwarded-Proto=https / X-Forwarded-For / Host —
# приложение это учитывает (uvicorn --proxy-headers). Положи файл в /etc/caddy/Caddyfile. # приложение это учитывает (uvicorn --proxy-headers). Положи файл в /etc/caddy/Caddyfile.
# Сертификаты — см. deploy/vps/README.md (fullchain = leaf + промежуточные одним файлом). # Сертификаты — см. deploy/vps/README.md (fullchain = leaf + промежуточные одним файлом).
#
# SSE (/api/events): отдельный handle БЕЗ encode и с flush_interval -1 — иначе сжатие/
# буферизация задержат поток событий (реал-тайм перестанет работать).
forbiddenstars.ru { forbiddenstars.ru {
tls /etc/caddy/certs/forbiddenstars.ru/fullchain.pem /etc/caddy/certs/forbiddenstars.ru/privkey.pem tls /etc/caddy/certs/forbiddenstars.ru/fullchain.pem /etc/caddy/certs/forbiddenstars.ru/privkey.pem
encode zstd gzip @sse path /api/events
reverse_proxy 127.0.0.1:9000 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 { forbidden-stars.ru {
tls /etc/caddy/certs/forbidden-stars.ru/fullchain.pem /etc/caddy/certs/forbidden-stars.ru/privkey.pem tls /etc/caddy/certs/forbidden-stars.ru/fullchain.pem /etc/caddy/certs/forbidden-stars.ru/privkey.pem
encode zstd gzip @sse path /api/events
reverse_proxy 127.0.0.1:9001 handle @sse {
reverse_proxy 127.0.0.1:9001 {
flush_interval -1
}
}
handle {
encode zstd gzip
reverse_proxy 127.0.0.1:9001
}
} }
+46 -1
View File
@@ -645,6 +645,23 @@ export interface paths {
patch?: never; patch?: never;
trace?: 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": { "/api/admin/auth/login": {
parameters: { parameters: {
query?: never; query?: never;
@@ -1507,6 +1524,8 @@ export interface components {
win_reason: "objectives" | "worlds" | "plastic" | "resources"; win_reason: "objectives" | "worlds" | "plastic" | "resources";
/** Overall Comment */ /** Overall Comment */
overall_comment?: string | null; overall_comment?: string | null;
/** Expected Version */
expected_version?: string | null;
}; };
/** MatchFinishParticipant */ /** MatchFinishParticipant */
MatchFinishParticipant: { MatchFinishParticipant: {
@@ -1627,6 +1646,8 @@ export interface components {
* @default false * @default false
*/ */
can_modify: boolean; can_modify: boolean;
/** Version */
version: string;
/** /**
* Participants * Participants
* @default [] * @default []
@@ -1648,6 +1669,8 @@ export interface components {
win_reason?: ("objectives" | "worlds" | "plastic" | "resources") | null; win_reason?: ("objectives" | "worlds" | "plastic" | "resources") | null;
/** Participants */ /** Participants */
participants?: components["schemas"]["ParticipantInput"][] | null; participants?: components["schemas"]["ParticipantInput"][] | null;
/** Expected Version */
expected_version?: string | null;
}; };
/** MeRead */ /** MeRead */
MeRead: { MeRead: {
@@ -2882,7 +2905,9 @@ export interface operations {
}; };
delete_match_api_matches__match_id__delete: { delete_match_api_matches__match_id__delete: {
parameters: { parameters: {
query?: never; query?: {
expected_version?: string | null;
};
header?: never; header?: never;
path: { path: {
match_id: number; 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: { admin_login_api_admin_auth_login_post: {
parameters: { parameters: {
query?: never; query?: never;
+5 -2
View File
@@ -4,8 +4,11 @@ export const queryClient = new QueryClient({
defaultOptions: { defaultOptions: {
queries: { queries: {
retry: false, retry: false,
refetchOnWindowFocus: false, // Свежесть: возврат во вкладку/восстановление сети → перезапрос (на случай, если
staleTime: 30_000, // SSE-соединение временно отвалилось). Точечные апдейты приходят пушем (useServerEvents).
refetchOnWindowFocus: true,
refetchOnReconnect: true,
staleTime: 15_000,
}, },
}, },
}); });
+4
View File
@@ -1,5 +1,7 @@
import { Outlet, useLocation, useNavigate } from "react-router-dom"; import { Outlet, useLocation, useNavigate } from "react-router-dom";
import { useMe } from "../hooks/auth";
import { useServerEvents } from "../hooks/useServerEvents";
import { BottomBar } from "./BottomBar"; import { BottomBar } from "./BottomBar";
const TITLES: Record<string, string> = { const TITLES: Record<string, string> = {
@@ -14,6 +16,8 @@ const TITLES: Record<string, string> = {
export function AppShell() { export function AppShell() {
const location = useLocation(); const location = useLocation();
const navigate = useNavigate(); const navigate = useNavigate();
const { data: me } = useMe();
useServerEvents(!!me); // живые обновления (SSE) только при наличии сессии
const title = const title =
TITLES[location.pathname] ?? TITLES[location.pathname] ??
(location.pathname.startsWith("/u/") ? "Профиль игрока" : "Forbidden Stars"); (location.pathname.startsWith("/u/") ? "Профиль игрока" : "Forbidden Stars");
+6 -2
View File
@@ -70,10 +70,14 @@ export function useFinishMatch() {
export function useDeleteMatch() { export function useDeleteMatch() {
const qc = useQueryClient(); const qc = useQueryClient();
return useMutation({ return useMutation({
mutationFn: async (matchId: number) => // expectedVersion — оптимистичная блокировка: отмена устаревшей версии вернёт 409 STALE_WRITE.
mutationFn: async (args: { matchId: number; expectedVersion?: string }) =>
unwrap( unwrap(
await api.DELETE("/api/matches/{match_id}", { 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(), onSuccess: () => qc.invalidateQueries(),
+56
View File
@@ -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]);
}
+24 -6
View File
@@ -26,7 +26,7 @@ interface FinishRow {
export function MatchDetailPage() { export function MatchDetailPage() {
const { matchId } = useParams(); const { matchId } = useParams();
const id = matchId ? Number(matchId) : null; const id = matchId ? Number(matchId) : null;
const { data: match, isLoading } = useMatch(id); const { data: match, isLoading, refetch } = useMatch(id);
const finish = useFinishMatch(); const finish = useFinishMatch();
const del = useDeleteMatch(); const del = useDeleteMatch();
const uploadAtt = useUploadMatchAttachment(id ?? 0); const uploadAtt = useUploadMatchAttachment(id ?? 0);
@@ -59,8 +59,11 @@ export function MatchDetailPage() {
const upd = (i: number, patch: Partial<FinishRow>) => const upd = (i: number, patch: Partial<FinishRow>) =>
setRows(finishRows.map((r, idx) => (idx === i ? { ...r, ...patch } : r))); setRows(finishRows.map((r, idx) => (idx === i ? { ...r, ...patch } : r)));
// Конфликт версий (кто-то изменил партию с другого устройства) → сообщаем и обновляем.
const isStale = (e: unknown) => e instanceof ApiError && e.code === "STALE_WRITE";
const submitFinish = async () => { const submitFinish = async () => {
if (!id) return; if (!id || !match) return;
setError(null); setError(null);
try { try {
await finish.mutateAsync({ await finish.mutateAsync({
@@ -73,19 +76,34 @@ export function MatchDetailPage() {
})), })),
win_reason: winReason, win_reason: winReason,
overall_comment: overall.trim() || null, overall_comment: overall.trim() || null,
expected_version: match.version,
}, },
}); });
toast.show("Партия завершена"); toast.show("Партия завершена");
} catch (e) { } catch (e) {
if (isStale(e)) {
toast.show("Партия изменилась на другом устройстве — обновлено");
refetch();
return;
}
setError(e instanceof ApiError ? e.message : "Не удалось завершить партию"); setError(e instanceof ApiError ? e.message : "Не удалось завершить партию");
} }
}; };
const remove = async () => { const remove = async () => {
if (!id) return; if (!id || !match) return;
await del.mutateAsync(id).catch(() => {}); try {
toast.show(inProgress ? "Партия отменена" : "Партия удалена"); await del.mutateAsync({ matchId: id, expectedVersion: match.version });
navigate(-1); 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( const sorted = [...match.participants].sort(