Уведомления: persistent + live (push-to-invalidate → pull)
Новая система уведомлений поверх готовой SSE-шины: события пишутся в БД,
живут 72 часа и чистятся (при чтении списка + фоновой задачей), всплывают
сверху экрана в момент прихода, доступны через колокольчик в правом верхнем
углу (бейдж непрочитанных + панель).
Типы: приглашение в группу (→ /group), старт/финиш партии участникам кроме
инициатора (→ /match/{id}). Титулы — готовый хелпер-задел (не подключён, т.к.
выдача титулов игрокам ещё не реализована).
Бэкенд: модель Notification + миграция 0008 (идемпотентная), notification_service,
notify.notifications_changed, роутер /api/notifications (GET + /read), триггеры
в groups/matches, фоновая чистка в lifespan, защита hub.publish от закрытого loop.
Фронт: useNotifications/useMarkNotificationsRead, NotificationBell/Panel/Toaster,
перекомпоновка top-bar, стили; useServerEvents знает тип notifications.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
@@ -37,11 +37,14 @@ class EventHub:
|
||||
def publish(self, user_ids: Iterable[int], event: dict[str, Any]) -> None:
|
||||
"""Доставить событие подписчикам с указанными user_id (потокобезопасно)."""
|
||||
loop = self._loop
|
||||
if loop is None:
|
||||
if loop is None or loop.is_closed():
|
||||
return
|
||||
ids = [uid for uid in user_ids if uid is not None]
|
||||
if ids:
|
||||
loop.call_soon_threadsafe(self._fanout, ids, event)
|
||||
try:
|
||||
loop.call_soon_threadsafe(self._fanout, ids, event)
|
||||
except RuntimeError:
|
||||
pass # loop закрылся между проверкой и вызовом (напр. при shutdown)
|
||||
|
||||
def _fanout(self, ids: list[int], event: dict[str, Any]) -> None:
|
||||
# Выполняется на потоке loop'а — доступ к _subs безопасен.
|
||||
|
||||
+33
-2
@@ -22,6 +22,7 @@ from app.routers import (
|
||||
groups,
|
||||
invitations,
|
||||
matches,
|
||||
notifications,
|
||||
reference,
|
||||
stats,
|
||||
users,
|
||||
@@ -71,6 +72,31 @@ class CSRFMiddleware:
|
||||
await self.app(scope, receive, send)
|
||||
|
||||
|
||||
_NOTIFICATIONS_PURGE_INTERVAL = 3600 # раз в час чистим протухшие уведомления (>72ч)
|
||||
|
||||
|
||||
async def _notifications_purge_loop() -> None:
|
||||
"""Фоновая чистка протухших уведомлений (single-worker безопасно). Чтобы удалялись
|
||||
«отовсюду» даже у неактивных пользователей (помимо очистки при чтении списка)."""
|
||||
import asyncio
|
||||
|
||||
from app.db.session import Session, engine
|
||||
from app.services import notification_service
|
||||
|
||||
def _purge_once() -> None:
|
||||
with Session(engine) as session:
|
||||
notification_service.purge_expired(session)
|
||||
|
||||
while True:
|
||||
try:
|
||||
await asyncio.sleep(_NOTIFICATIONS_PURGE_INTERVAL)
|
||||
await asyncio.to_thread(_purge_once)
|
||||
except asyncio.CancelledError:
|
||||
break
|
||||
except Exception as exc: # noqa: BLE001
|
||||
logging.getLogger("fs").warning("Чистка уведомлений пропущена: %s", exc)
|
||||
|
||||
|
||||
@asynccontextmanager
|
||||
async def _lifespan(_app: FastAPI):
|
||||
# SSE-шина публикует из sync-роутеров в этот event-loop — сохраняем ссылку (все окружения).
|
||||
@@ -91,7 +117,12 @@ async def _lifespan(_app: FastAPI):
|
||||
logging.getLogger("fs").warning(
|
||||
"Стартовый bootstrap пропущен (примените миграции): %s", exc
|
||||
)
|
||||
yield
|
||||
|
||||
purge_task = asyncio.create_task(_notifications_purge_loop())
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
purge_task.cancel()
|
||||
|
||||
|
||||
def create_app() -> FastAPI:
|
||||
@@ -135,7 +166,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,
|
||||
events.router, admin.router]
|
||||
events.router, notifications.router, admin.router]
|
||||
for r in api_routers:
|
||||
app.include_router(r, prefix="/api")
|
||||
|
||||
|
||||
@@ -336,6 +336,33 @@ class MatchAttachment(SQLModel, table=True):
|
||||
created_at: datetime = Field(default_factory=_utcnow, nullable=False)
|
||||
|
||||
|
||||
# ─── Уведомления ─────────────────────────────────────────────────────────────
|
||||
|
||||
class Notification(SQLModel, table=True):
|
||||
"""Персистентное уведомление игроку (приглашение, старт/финиш партии, титул и т.д.).
|
||||
|
||||
Текст (`title`/`body`, RU) и ссылку (`link` — относительный SPA-путь) рендерит сервер —
|
||||
фронт лишь отображает. Хранятся 72 часа; протухшие чистятся фоном и при чтении списка."""
|
||||
|
||||
__tablename__ = "notifications"
|
||||
__table_args__ = (
|
||||
Index("ix_notifications_user_created", "user_id", "created_at"),
|
||||
)
|
||||
|
||||
id: int | None = Field(default=None, primary_key=True)
|
||||
user_id: int = Field(
|
||||
sa_column=Column(
|
||||
Integer, ForeignKey("users.id", ondelete="CASCADE"), nullable=False, index=True
|
||||
)
|
||||
)
|
||||
type: str = Field(sa_column=Column(String(32), nullable=False))
|
||||
title: str = Field(sa_column=Column(String(255), nullable=False))
|
||||
body: str | None = Field(default=None, sa_column=Column(Text, nullable=True))
|
||||
link: str | None = Field(default=None, sa_column=Column(String(255), nullable=True))
|
||||
read_at: datetime | None = Field(default=None, sa_column=Column(DateTime, nullable=True))
|
||||
created_at: datetime = Field(default_factory=_utcnow, nullable=False, index=True)
|
||||
|
||||
|
||||
# ─── Журнал аудита ───────────────────────────────────────────────────────────
|
||||
|
||||
class AuditLog(SQLModel, table=True):
|
||||
|
||||
@@ -14,6 +14,7 @@ from app.services import (
|
||||
group_service,
|
||||
invitation_service,
|
||||
membership_service,
|
||||
notification_service,
|
||||
notify,
|
||||
stats_service,
|
||||
)
|
||||
@@ -155,6 +156,9 @@ def invite_member(
|
||||
)
|
||||
session.commit()
|
||||
notify.invitations_changed(invited.id) # type: ignore[arg-type] # живое появление у получателя
|
||||
notification_service.invited_to_group(
|
||||
session, invited.id, group.name, user.nickname # type: ignore[arg-type]
|
||||
)
|
||||
return s.InvitationRead(
|
||||
id=inv.id, # type: ignore[arg-type]
|
||||
group_id=group_id,
|
||||
|
||||
@@ -16,6 +16,7 @@ from app.services import (
|
||||
audit_service,
|
||||
group_service,
|
||||
match_service,
|
||||
notification_service,
|
||||
notify,
|
||||
user_service,
|
||||
)
|
||||
@@ -121,6 +122,7 @@ def start_match(
|
||||
)
|
||||
session.commit()
|
||||
notify.match_changed(session, match)
|
||||
notification_service.match_started(session, match, user.id) # type: ignore[arg-type]
|
||||
return build_match_read(session, match, can_modify=match_service.can_modify(session, match, user))
|
||||
|
||||
|
||||
@@ -160,6 +162,7 @@ def finish_match(
|
||||
)
|
||||
session.commit()
|
||||
notify.match_changed(session, match)
|
||||
notification_service.match_finished(session, match, user.id) # type: ignore[arg-type]
|
||||
return build_match_read(session, match, can_modify=match_service.can_modify(session, match, user))
|
||||
|
||||
|
||||
|
||||
@@ -0,0 +1,34 @@
|
||||
"""Уведомления игрока: список (pull) и отметка прочитанными.
|
||||
|
||||
Появление в реальном времени обеспечивает SSE-сигнал `{type:"notifications"}` — по нему клиент
|
||||
перезапрашивает этот список. Хранение — 72 часа; протухшие чистятся при чтении (и фоном)."""
|
||||
from __future__ import annotations
|
||||
|
||||
from fastapi import APIRouter, Depends
|
||||
from sqlmodel import Session
|
||||
|
||||
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 notification_service
|
||||
|
||||
router = APIRouter(prefix="/notifications", tags=["notifications"])
|
||||
|
||||
|
||||
@router.get("", response_model=s.NotificationList)
|
||||
def my_notifications(
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(get_current_user),
|
||||
) -> dict:
|
||||
return notification_service.list_for_user(session, user.id) # type: ignore[arg-type]
|
||||
|
||||
|
||||
@router.post("/read", response_model=s.OkResponse)
|
||||
def mark_read(
|
||||
body: s.NotificationMarkRead,
|
||||
session: Session = Depends(get_session),
|
||||
user: User = Depends(get_current_user),
|
||||
) -> s.OkResponse:
|
||||
notification_service.mark_read(session, user.id, body.ids) # type: ignore[arg-type]
|
||||
return s.OkResponse()
|
||||
@@ -157,6 +157,27 @@ class InvitationRead(BaseModel):
|
||||
created_at: str
|
||||
|
||||
|
||||
# ─── Уведомления ───────────────────────────────────────────────────────────────
|
||||
|
||||
class NotificationRead(BaseModel):
|
||||
id: int
|
||||
type: str
|
||||
title: str
|
||||
body: str | None = None
|
||||
link: str | None = None
|
||||
read_at: str | None = None
|
||||
created_at: str
|
||||
|
||||
|
||||
class NotificationList(BaseModel):
|
||||
items: list[NotificationRead]
|
||||
unread_count: int
|
||||
|
||||
|
||||
class NotificationMarkRead(BaseModel):
|
||||
ids: list[int] | None = None
|
||||
|
||||
|
||||
# ─── Партии ──────────────────────────────────────────────────────────────────
|
||||
|
||||
class RandomizeRequest(BaseModel):
|
||||
|
||||
@@ -0,0 +1,190 @@
|
||||
"""Персистентные уведомления игрокам (приглашения, старт/финиш партии, титулы).
|
||||
|
||||
Запись в БД + живой сигнал по SSE (`notify.notifications_changed`) — клиент по сигналу
|
||||
подтягивает список (`GET /api/notifications`). Текст (RU) и ссылку рендерим здесь, на сервере.
|
||||
Хранение — 72 часа; протухшие удаляются при чтении списка и фоновой задачей.
|
||||
"""
|
||||
from __future__ import annotations
|
||||
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
from sqlmodel import Session, select
|
||||
|
||||
from app.core.timeutil import iso_utc
|
||||
from app.models import Group, Match, MatchParticipant, Notification
|
||||
from app.services import notify
|
||||
|
||||
RETENTION_HOURS = 72
|
||||
|
||||
|
||||
def _now() -> datetime:
|
||||
return datetime.now(timezone.utc)
|
||||
|
||||
|
||||
def _cutoff() -> datetime:
|
||||
# Наивный UTC — сопоставимо с тем, как SQLite отдаёт сохранённые datetime.
|
||||
return _now().replace(tzinfo=None) - timedelta(hours=RETENTION_HOURS)
|
||||
|
||||
|
||||
# ─── Базовые операции ─────────────────────────────────────────────────────────
|
||||
|
||||
def create_for(
|
||||
session: Session,
|
||||
user_ids,
|
||||
*,
|
||||
type: str,
|
||||
title: str,
|
||||
body: str | None = None,
|
||||
link: str | None = None,
|
||||
) -> list[Notification]:
|
||||
"""Создать одно уведомление каждому из user_ids (один commit) и пушнуть им SSE-сигнал."""
|
||||
ids = [uid for uid in dict.fromkeys(user_ids) if uid is not None]
|
||||
if not ids:
|
||||
return []
|
||||
rows = [
|
||||
Notification(user_id=uid, type=type, title=title, body=body, link=link) for uid in ids
|
||||
]
|
||||
session.add_all(rows)
|
||||
session.commit()
|
||||
for uid in ids:
|
||||
notify.notifications_changed(uid)
|
||||
return rows
|
||||
|
||||
|
||||
def create(
|
||||
session: Session,
|
||||
user_id: int,
|
||||
*,
|
||||
type: str,
|
||||
title: str,
|
||||
body: str | None = None,
|
||||
link: str | None = None,
|
||||
) -> Notification | None:
|
||||
rows = create_for(session, [user_id], type=type, title=title, body=body, link=link)
|
||||
return rows[0] if rows else None
|
||||
|
||||
|
||||
def purge_expired(session: Session) -> int:
|
||||
"""Удалить уведомления старше RETENTION_HOURS. Возвращает число удалённых."""
|
||||
rows = session.exec(
|
||||
select(Notification).where(Notification.created_at < _cutoff())
|
||||
).all()
|
||||
for row in rows:
|
||||
session.delete(row)
|
||||
if rows:
|
||||
session.commit()
|
||||
return len(rows)
|
||||
|
||||
|
||||
def list_for_user(session: Session, user_id: int) -> dict:
|
||||
"""Свежие (<72ч) уведомления пользователя + число непрочитанных. Чистит протухшие."""
|
||||
purge_expired(session)
|
||||
rows = session.exec(
|
||||
select(Notification)
|
||||
.where(Notification.user_id == user_id)
|
||||
.order_by(Notification.created_at.desc())
|
||||
).all()
|
||||
items = [
|
||||
{
|
||||
"id": n.id,
|
||||
"type": n.type,
|
||||
"title": n.title,
|
||||
"body": n.body,
|
||||
"link": n.link,
|
||||
"read_at": iso_utc(n.read_at),
|
||||
"created_at": iso_utc(n.created_at),
|
||||
}
|
||||
for n in rows
|
||||
]
|
||||
unread = sum(1 for n in rows if n.read_at is None)
|
||||
return {"items": items, "unread_count": unread}
|
||||
|
||||
|
||||
def mark_read(session: Session, user_id: int, ids: list[int] | None = None) -> int:
|
||||
"""Отметить прочитанными все непрочитанные пользователя (или указанные по id)."""
|
||||
query = select(Notification).where(
|
||||
Notification.user_id == user_id, Notification.read_at.is_(None)
|
||||
)
|
||||
if ids:
|
||||
query = query.where(Notification.id.in_(ids))
|
||||
rows = session.exec(query).all()
|
||||
now = _now()
|
||||
for row in rows:
|
||||
row.read_at = now
|
||||
session.add(row)
|
||||
if rows:
|
||||
session.commit()
|
||||
return len(rows)
|
||||
|
||||
|
||||
# ─── Типовые хелперы под триггеры ──────────────────────────────────────────────
|
||||
|
||||
def _match_participant_ids(session: Session, match_id: int) -> list[int]:
|
||||
return list(
|
||||
session.exec(
|
||||
select(MatchParticipant.user_id).where(MatchParticipant.match_id == match_id)
|
||||
).all()
|
||||
)
|
||||
|
||||
|
||||
def _group_name(session: Session, group_id: int) -> str | None:
|
||||
group = session.get(Group, group_id)
|
||||
return group.name if group else None
|
||||
|
||||
|
||||
def invited_to_group(
|
||||
session: Session, user_id: int, group_name: str, inviter_nickname: str | None = None
|
||||
) -> None:
|
||||
create(
|
||||
session,
|
||||
user_id,
|
||||
type="invite",
|
||||
title=f"Приглашение в группу «{group_name}»",
|
||||
body=f"пригласил: {inviter_nickname}" if inviter_nickname else None,
|
||||
link="/group",
|
||||
)
|
||||
|
||||
|
||||
def match_started(session: Session, match: Match, actor_id: int) -> None:
|
||||
ids = [uid for uid in _match_participant_ids(session, match.id) if uid != actor_id] # type: ignore[arg-type]
|
||||
name = _group_name(session, match.group_id)
|
||||
create_for(
|
||||
session,
|
||||
ids,
|
||||
type="match_started",
|
||||
title="Началась партия",
|
||||
body=f"Группа «{name}»" if name else None,
|
||||
link=f"/match/{match.id}",
|
||||
)
|
||||
|
||||
|
||||
def match_finished(session: Session, match: Match, actor_id: int) -> None:
|
||||
ids = [uid for uid in _match_participant_ids(session, match.id) if uid != actor_id] # type: ignore[arg-type]
|
||||
name = _group_name(session, match.group_id)
|
||||
create_for(
|
||||
session,
|
||||
ids,
|
||||
type="match_finished",
|
||||
title="Партия завершена",
|
||||
body=f"Группа «{name}»" if name else None,
|
||||
link=f"/match/{match.id}",
|
||||
)
|
||||
|
||||
|
||||
def title_earned(session: Session, user_id: int, achievement_slug: str) -> None:
|
||||
"""ЗАДЕЛ: уведомление о полученном титуле. Пока нигде не вызывается — заработает,
|
||||
когда появится механизм выдачи титулов игрокам."""
|
||||
name = achievement_slug
|
||||
try:
|
||||
from app.services import achievement_service
|
||||
|
||||
name = achievement_service.get(achievement_slug).get("name", achievement_slug)
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
create(
|
||||
session,
|
||||
user_id,
|
||||
type="title",
|
||||
title=f"Вы получили титул: {name}",
|
||||
link="/account",
|
||||
)
|
||||
@@ -44,3 +44,8 @@ def group_changed(session: Session, group_id: int, extra_user_ids: list[int] | N
|
||||
def invitations_changed(user_id: int) -> None:
|
||||
"""У пользователя изменился список приглашений."""
|
||||
hub.publish([user_id], {"type": "invitations"})
|
||||
|
||||
|
||||
def notifications_changed(user_id: int) -> None:
|
||||
"""У пользователя появилось/изменилось уведомление — пусть подтянет список."""
|
||||
hub.publish([user_id], {"type": "notifications"})
|
||||
|
||||
Reference in New Issue
Block a user