From ed30a21842953f39a903d80fc857c7ca5185011d Mon Sep 17 00:00:00 2001 From: Vladless Date: Sun, 19 Apr 2026 20:59:15 +0000 Subject: [PATCH] auth sessions + flow schema fixes --- api/depends.py | 18 +++- api/main.py | 6 ++ api/v2/routes/auth/_common.py | 8 ++ api/v2/routes/auth/google.py | 2 +- api/v2/routes/auth/password.py | 20 ++-- api/v2/routes/auth/session.py | 82 +++++++++++++- api/v2/routes/auth/telegram.py | 14 +-- api/v2/routes/auth/yandex.py | 2 +- api/v2/routes/flows.py | 10 +- api/v2/routes/web.py | 53 ++++++++- api/v2/schemas/flows.py | 31 +++++- api/v2/schemas/identities.py | 15 +++ api/v2/schemas/web.py | 17 +++ database/__init__.py | 1 + database/identities.py | 75 +++++++++---- database/identity_sessions.py | 149 ++++++++++++++++++++++++++ database/migrations/schema_upgrade.py | 53 +++++++++ database/models/__init__.py | 2 + database/models/identity_session.py | 39 +++++++ utils/versioning.py | 2 +- 20 files changed, 545 insertions(+), 54 deletions(-) create mode 100644 database/identity_sessions.py create mode 100644 database/models/identity_session.py diff --git a/api/depends.py b/api/depends.py index 2d83c9c4..05bd6f59 100644 --- a/api/depends.py +++ b/api/depends.py @@ -7,10 +7,13 @@ from fastapi import Depends, HTTPException, Header, Query, Request, Response from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession +from datetime import datetime + from audit import set_api_actor from database import ( async_session_maker, identities as idb, + identity_sessions as idsess, ) from database.access.resolution import ResolvedActor, resolve_actor_from_identity from database.models import Admin, Identity @@ -148,11 +151,20 @@ async def _identity_from_cookie(session: AsyncSession, request: Request | None) if not token: return None token_hash = hash_token(token) - identity = await idb.get_identity_by_token_hash(session, token_hash) + sess = await idsess.get_session_by_token_hash(session, token_hash) + if sess is None: + return None + if sess.expires_at is not None and sess.expires_at <= datetime.utcnow(): + return None + identity = await idb.get_identity_by_id(session, sess.identity_id) if identity is None: return None - if idb._is_token_expired(identity): - return None + await idsess.touch_session_last_seen(session, sess) + if request is not None: + try: + request.state.auth_session = sess + except Exception: + pass return identity diff --git a/api/main.py b/api/main.py index 97706fc0..be58dab6 100644 --- a/api/main.py +++ b/api/main.py @@ -53,6 +53,12 @@ async def security_and_cache_middleware(request: Request, call_next): response.headers.setdefault("X-XSS-Protection", "1; mode=block") content_type = response.headers.get("content-type", "") + path = request.url.path + + if path.startswith("/api/web/uploads/") and request.method == "GET" and response.status_code == 200: + response.headers.setdefault("Cache-Control", "public, max-age=3600, stale-while-revalidate=86400") + return response + if request.method == "GET" and response.status_code == 200 and "application/json" in content_type: body = b"" async for chunk in response.body_iterator: diff --git a/api/v2/routes/auth/_common.py b/api/v2/routes/auth/_common.py index 9ba06c8e..30925f89 100644 --- a/api/v2/routes/auth/_common.py +++ b/api/v2/routes/auth/_common.py @@ -2,6 +2,7 @@ from fastapi import Request from sqlalchemy import text from sqlalchemy.ext.asyncio import AsyncSession +from api.v2.schemas.identities import IdentityResponse, LoginResponse from config import API_TOKEN_TTL_DAYS from logger import logger from utils.referral_codes import encode_partner_code @@ -10,6 +11,13 @@ from utils.referral_codes import encode_partner_code TOKEN_TTL_HINT = "бессрочно" if API_TOKEN_TTL_DAYS is None else f"{API_TOKEN_TTL_DAYS} дн." TELEGRAM_LOGIN_MAX_AGE = 86400 + +def build_login_response(identity) -> LoginResponse: + return LoginResponse( + identity_id=identity.id, + identity=IdentityResponse.model_validate(identity), + ) + _TRUSTED_PROXY_CIDRS: list[str] = [] diff --git a/api/v2/routes/auth/google.py b/api/v2/routes/auth/google.py index 666d8697..ae6b0222 100644 --- a/api/v2/routes/auth/google.py +++ b/api/v2/routes/auth/google.py @@ -191,7 +191,7 @@ async def google_callback( email=email if (email and email_verified) else None, ) await bind_identity_actor(request, session, identity) - token = await idb.issue_token_for_identity(session, identity) + token = await idb.issue_token_for_identity(session, identity, request=request) logger.info( "[Auth] Login success: identity={}, google_sub={}, ip={}, method=google", identity.id, diff --git a/api/v2/routes/auth/password.py b/api/v2/routes/auth/password.py index 1668169d..0684adb6 100644 --- a/api/v2/routes/auth/password.py +++ b/api/v2/routes/auth/password.py @@ -9,7 +9,7 @@ from api.depends import ( set_auth_cookie, set_is_admin_cookie, ) -from api.v2.routes.auth._common import TOKEN_TTL_HINT, _client_ip +from api.v2.routes.auth._common import TOKEN_TTL_HINT, _client_ip, build_login_response from api.v2.schemas.identities import ( ConfirmPasswordResetRequest, LoginByCodeRequest, @@ -114,7 +114,9 @@ async def register_by_email( referrer_user = await resolve_user_optional(session, referrer_legacy) if referrer_user is None: raise HTTPException(status_code=400, detail="Код приглашения недействителен") - identity, token = await idb.create_identity_with_token(session, email=email, password=body.password) + identity, token = await idb.create_identity_with_token( + session, email=email, password=body.password, request=request + ) await bind_identity_actor(request, session, identity) billing_user_id = await idb.ensure_billing_user_for_identity(session, identity) if referrer_user is not None and not await get_referral_by_referred_id(session, billing_user_id): @@ -176,14 +178,14 @@ async def login( raise except Exception as e: logger.warning("[Auth] Ошибка rate-limit проверки для email-логина: {}", e) - result = await idb.login_by_email(session, email, body.password) + result = await idb.login_by_email(session, email, body.password, request=request) if not result: from database.setup.web_admin_bootstrap import ensure_web_admin try: await ensure_web_admin(session) await session.flush() - result = await idb.login_by_email(session, email, body.password) + result = await idb.login_by_email(session, email, body.password, request=request) except Exception as exc: logger.warning("[Auth] lazy web-admin bootstrap failed: {}", exc) if not result: @@ -212,7 +214,7 @@ async def login( logger.info("[Auth] Login success: identity={}, email={}, ip={}, method=password", identity.id, email, ip) set_auth_cookie(response, token, request) set_is_admin_cookie(response, identity, request) - return LoginResponse(identity_id=identity.id) + return build_login_response(identity) @router.post("/send-login-code") @@ -329,7 +331,7 @@ async def login_by_code( sa_update(IdentityModel).where(IdentityModel.id == identity.id).values(email_verified=True) ) await bind_identity_actor(request, session, identity) - token = await idb.issue_token_for_identity(session, identity) + token = await idb.issue_token_for_identity(session, identity, request=request) if getattr(identity, "is_admin", False): from database.site_state import mark_site_initialized @@ -337,7 +339,7 @@ async def login_by_code( logger.info("[Auth] Login success: identity={}, email={}, method=code", identity.id, email_norm) set_auth_cookie(response, token, request) set_is_admin_cookie(response, identity, request) - return LoginResponse(identity_id=identity.id) + return build_login_response(identity) @router.post("/request-password-reset") @@ -422,7 +424,7 @@ async def confirm_password_reset( if not updated: raise HTTPException(status_code=400, detail="Не удалось обновить пароль") await bind_identity_actor(request, session, updated) - token = await idb.issue_token_for_identity(session, updated) + token = await idb.issue_token_for_identity(session, updated, request=request) set_auth_cookie(response, token, request) set_is_admin_cookie(response, updated, request) - return LoginResponse(identity_id=updated.id) + return build_login_response(updated) diff --git a/api/v2/routes/auth/session.py b/api/v2/routes/auth/session.py index b0f40376..1104bd6b 100644 --- a/api/v2/routes/auth/session.py +++ b/api/v2/routes/auth/session.py @@ -7,12 +7,16 @@ from api.depends import ( clear_auth_cookie, get_request_actor, get_session, + hash_token, verify_identity_token, ) +from api.depends import AUTH_COOKIE_NAME from api.v2.routes.auth._common import _resolve_partner_snapshot from api.v2.schemas.identities import ( ChangePasswordRequest, IdentityResponse, + IdentitySessionItem, + IdentitySessionsResponse, SetPasswordRequest, ) from api.v2.schemas.web_public import AccountSummaryResponse @@ -21,6 +25,7 @@ from database import ( get_keys, get_trial, identities as idb, + identity_sessions as idsess, ) from database.models import CouponUsage, Gift, GiftUsage from database.referrals import get_referral_stats @@ -39,12 +44,87 @@ async def me( return IdentityResponse.model_validate(identity) +def _current_token_hash(request: Request) -> str | None: + raw = request.cookies.get(AUTH_COOKIE_NAME) + if not raw or not raw.strip(): + return None + return hash_token(raw.strip()) + + +@router.get("/sessions", response_model=IdentitySessionsResponse) +async def list_my_sessions( + request: Request, + session: AsyncSession = Depends(get_session), + identity=Depends(verify_identity_token), +): + """Возвращает активные сессии текущей identity (все устройства).""" + current_hash = _current_token_hash(request) + rows = await idsess.list_sessions_for_identity(session, identity.id) + items = [ + IdentitySessionItem( + id=row.id, + device_label=row.device_label, + ip=row.ip, + created_at=row.created_at, + last_seen_at=row.last_seen_at, + expires_at=row.expires_at, + is_current=bool(current_hash and row.token_hash == current_hash), + ) + for row in rows + ] + return IdentitySessionsResponse(sessions=items) + + +@router.delete("/sessions/{session_id}") +async def revoke_my_session( + session_id: str, + request: Request, + response: Response, + session: AsyncSession = Depends(get_session), + identity=Depends(verify_identity_token), +): + """Удаляет одну сессию текущей identity. Если удалена текущая — очищаем cookie.""" + ok = await idsess.delete_session_by_id( + session, session_id=session_id, identity_id=identity.id + ) + if not ok: + raise HTTPException(status_code=404, detail="Сессия не найдена") + current_hash = _current_token_hash(request) + rows = await idsess.list_sessions_for_identity(session, identity.id) + if current_hash and not any(r.token_hash == current_hash for r in rows): + clear_auth_cookie(response, request) + return {"ok": True} + + +@router.post("/sessions/revoke-others") +async def revoke_other_sessions( + request: Request, + session: AsyncSession = Depends(get_session), + identity=Depends(verify_identity_token), +): + """Удаляет все сессии текущей identity кроме текущей.""" + current_hash = _current_token_hash(request) + if not current_hash: + raise HTTPException(status_code=400, detail="Текущая сессия не определена") + removed = await idsess.delete_other_sessions( + session, identity_id=identity.id, keep_token_hash=current_hash + ) + return {"ok": True, "removed": removed} + + @router.post("/logout") async def logout( request: Request, response: Response, + session: AsyncSession = Depends(get_session), ): - """Очищает auth cookie. Не требует валидной сессии — всегда возвращает ok.""" + """Удаляет текущую сессию из БД и очищает auth cookie. Не требует валидной сессии.""" + raw = request.cookies.get(AUTH_COOKIE_NAME) + if raw and raw.strip(): + try: + await idsess.delete_session_by_token_hash(session, hash_token(raw.strip())) + except Exception: + pass clear_auth_cookie(response, request) return {"ok": True} diff --git a/api/v2/routes/auth/telegram.py b/api/v2/routes/auth/telegram.py index 9295302a..ba503d90 100644 --- a/api/v2/routes/auth/telegram.py +++ b/api/v2/routes/auth/telegram.py @@ -12,7 +12,7 @@ from api.depends import ( set_is_admin_cookie, verify_identity_token, ) -from api.v2.routes.auth._common import TELEGRAM_LOGIN_MAX_AGE, TOKEN_TTL_HINT, _client_ip +from api.v2.routes.auth._common import TELEGRAM_LOGIN_MAX_AGE, TOKEN_TTL_HINT, _client_ip, build_login_response from api.v2.schemas.identities import ( IdentityResponse, LinkTelegramRequest, @@ -61,7 +61,7 @@ async def login_telegram( raise HTTPException(status_code=401, detail="Неверная подпись или устаревшие данные от Telegram") identity = await idb.get_or_create_identity_for_tg(session, body.id) await bind_identity_actor(request, session, identity) - token = await idb.issue_token_for_identity(session, identity) + token = await idb.issue_token_for_identity(session, identity, request=request) logger.info( "[Auth] Login success: identity={}, tg_id={}, ip={}, method=telegram_widget", identity.id, @@ -70,7 +70,7 @@ async def login_telegram( ) set_auth_cookie(response, token, request) set_is_admin_cookie(response, identity, request) - return LoginResponse(identity_id=identity.id) + return build_login_response(identity) @router.post("/login-telegram-webapp", response_model=LoginResponse) @@ -91,7 +91,7 @@ async def login_telegram_webapp( raise HTTPException(status_code=401, detail="Не удалось определить пользователя из initData") identity = await idb.get_or_create_identity_for_tg(session, int(tg_id)) await bind_identity_actor(request, session, identity) - token = await idb.issue_token_for_identity(session, identity) + token = await idb.issue_token_for_identity(session, identity, request=request) logger.info( "[Auth] Login success: identity={}, tg_id={}, ip={}, method=telegram_webapp", identity.id, @@ -100,7 +100,7 @@ async def login_telegram_webapp( ) set_auth_cookie(response, token, request) set_is_admin_cookie(response, identity, request) - return LoginResponse(identity_id=identity.id) + return build_login_response(identity) class LoginTelegramOIDCRequest(BaseModel): @@ -192,7 +192,7 @@ async def login_telegram_oidc( identity = await idb.get_or_create_identity_for_tg(session, tg_id_int) await bind_identity_actor(request, session, identity) - token = await idb.issue_token_for_identity(session, identity) + token = await idb.issue_token_for_identity(session, identity, request=request) if getattr(identity, "is_admin", False): from database.site_state import mark_site_initialized @@ -206,7 +206,7 @@ async def login_telegram_oidc( ) set_auth_cookie(response, token, request) set_is_admin_cookie(response, identity, request) - return LoginResponse(identity_id=identity.id) + return build_login_response(identity) @router.post("/link-telegram", response_model=IdentityResponse) diff --git a/api/v2/routes/auth/yandex.py b/api/v2/routes/auth/yandex.py index 5fb45e4d..8f64fea6 100644 --- a/api/v2/routes/auth/yandex.py +++ b/api/v2/routes/auth/yandex.py @@ -188,7 +188,7 @@ async def yandex_callback( email=email, ) await bind_identity_actor(request, session, identity) - token = await idb.issue_token_for_identity(session, identity) + token = await idb.issue_token_for_identity(session, identity, request=request) logger.info( "[Auth] Login success: identity={}, yandex_sub={}, ip={}, method=yandex", identity.id, diff --git a/api/v2/routes/flows.py b/api/v2/routes/flows.py index f65ad63c..059cd3d7 100644 --- a/api/v2/routes/flows.py +++ b/api/v2/routes/flows.py @@ -26,7 +26,7 @@ def _flow_to_response(flow: WebFlow) -> FlowResponse: ) -@router.get("/flows/{flow_id}", response_model=FlowResponse) +@router.get("/flows/{flow_id}", response_model=FlowResponse, response_model_by_alias=True) async def get_flow_public(flow_id: str, session: AsyncSession = Depends(get_session)): flow = await session.get(WebFlow, flow_id) if not flow: @@ -34,7 +34,7 @@ async def get_flow_public(flow_id: str, session: AsyncSession = Depends(get_sess return _flow_to_response(flow) -@router.get("/admin/flows", response_model=list[FlowResponse]) +@router.get("/admin/flows", response_model=list[FlowResponse], response_model_by_alias=True) async def list_flows( session: AsyncSession = Depends(get_session), _identity=Depends(verify_identity_admin), @@ -43,7 +43,7 @@ async def list_flows( return [_flow_to_response(f) for f in result.scalars().all()] -@router.get("/admin/flows/{flow_id}", response_model=FlowResponse) +@router.get("/admin/flows/{flow_id}", response_model=FlowResponse, response_model_by_alias=True) async def get_flow_admin( flow_id: str, session: AsyncSession = Depends(get_session), @@ -55,7 +55,7 @@ async def get_flow_admin( return _flow_to_response(flow) -@router.post("/admin/flows", response_model=FlowResponse, status_code=201) +@router.post("/admin/flows", response_model=FlowResponse, status_code=201, response_model_by_alias=True) async def create_flow( body: FlowCreate, session: AsyncSession = Depends(get_session), @@ -81,7 +81,7 @@ async def create_flow( return _flow_to_response(flow) -@router.put("/admin/flows/{flow_id}", response_model=FlowResponse) +@router.put("/admin/flows/{flow_id}", response_model=FlowResponse, response_model_by_alias=True) async def update_flow( flow_id: str, body: FlowUpdate, diff --git a/api/v2/routes/web.py b/api/v2/routes/web.py index 644b46e9..ffccfd54 100644 --- a/api/v2/routes/web.py +++ b/api/v2/routes/web.py @@ -14,6 +14,9 @@ from api.depends import get_session, verify_identity_admin from database.site_revision import bump_site_revision from api.v2.schemas import WebBlockResponse, WebPageResponse, WebPageUpdate, WebTheme from api.v2.schemas.web import ( + WebPageSaveResponse, + WebPageThemeResponse, + WebPageThemeUpdate, WebPageVariantCreate, WebPageVariantSummary, WebPageVariantUpdate, @@ -282,11 +285,51 @@ async def get_web_page( return await _build_page_response(session, slug, current, variants) -@router.put("/api/web/pages/{slug}", response_model=WebPageResponse) +@router.get("/api/web/pages/{slug}/theme", response_model=WebPageThemeResponse) +async def get_web_page_theme( + slug: str, + variant: str | None = Query(default=None), + session: AsyncSession = Depends(get_session), +): + """Возвращает только theme tokens страницы (без блоков/вариантов). Используется для fallback-темы cabinet-страниц.""" + if not slug or len(slug) > 64 or not _SLUG_RE.match(slug): + raise HTTPException(400, "Некорректный slug страницы") + current, _ = await _resolve_variant(session, slug, variant) + return WebPageThemeResponse( + slug=slug, + variant_key=current.variant_key, + tokens=dict(current.theme_tokens or {}), + ) + + +@router.put("/api/web/pages/{slug}/theme", response_model=WebPageThemeResponse) +async def update_web_page_theme( + slug: str, + body: WebPageThemeUpdate, + variant: str | None = Query(default=None), + session: AsyncSession = Depends(get_session), + identity=Depends(verify_identity_admin), +): + """Обновляет только theme_tokens страницы (без трогания блоков). Используется для sync темы между страницами.""" + if not slug or len(slug) > 64 or not _SLUG_RE.match(slug): + raise HTTPException(400, "Некорректный slug страницы") + current, _ = await _resolve_variant(session, slug, variant) + current.theme_tokens = body.tokens + await session.flush() + await bump_site_revision(session) + return WebPageThemeResponse( + slug=slug, + variant_key=current.variant_key, + tokens=dict(current.theme_tokens or {}), + ) + + +@router.put("/api/web/pages/{slug}") async def update_web_page( slug: str, body: WebPageUpdate, variant: str | None = Query(default=None), + minimal: bool = Query(default=False), session: AsyncSession = Depends(get_session), identity=Depends(verify_identity_admin), ): @@ -310,6 +353,14 @@ async def update_web_page( await bump_site_revision(session) refreshed_variants = await _list_variants(session, slug) refreshed_current = next((item for item in refreshed_variants if item.id == current.id), current) + if minimal: + active = next((item for item in refreshed_variants if item.is_active), refreshed_current) + return WebPageSaveResponse( + slug=slug, + variant_key=refreshed_current.variant_key, + active_variant_key=active.variant_key, + variants=[_variant_summary(item) for item in refreshed_variants], + ) return await _build_page_response(session, slug, refreshed_current, refreshed_variants) diff --git a/api/v2/schemas/flows.py b/api/v2/schemas/flows.py index fadbc528..a41754bc 100644 --- a/api/v2/schemas/flows.py +++ b/api/v2/schemas/flows.py @@ -1,37 +1,60 @@ from __future__ import annotations -from typing import Any +from typing import Any, Literal -from pydantic import BaseModel +from pydantic import BaseModel, ConfigDict +from pydantic.alias_generators import to_camel + + +_FLOW_CONFIG = ConfigDict(populate_by_name=True, alias_generator=to_camel) class EdgeConditionSchema(BaseModel): + model_config = _FLOW_CONFIG + field: str operator: str value: Any = None +class EdgeConditionGroupSchema(BaseModel): + model_config = _FLOW_CONFIG + + logic: Literal["and", "or"] = "and" + conditions: list[EdgeConditionSchema] = [] + + class FlowEdgeSchema(BaseModel): + model_config = _FLOW_CONFIG + id: str source: str target: str condition: EdgeConditionSchema | None = None + condition_group: EdgeConditionGroupSchema | None = None label: str | None = None priority: int | None = None class FlowNodeSchema(BaseModel): + model_config = _FLOW_CONFIG + id: str type: str label: str label_en: str | None = None enabled: bool = True page_slug: str | None = None + cabinet_tab: str | None = None + screen_group: str | None = None + screen_id: str | None = None config: dict = {} position: dict class FlowResponse(BaseModel): + model_config = _FLOW_CONFIG + id: str name: str nodes: list[FlowNodeSchema] @@ -41,6 +64,8 @@ class FlowResponse(BaseModel): class FlowUpdate(BaseModel): + model_config = _FLOW_CONFIG + name: str | None = None nodes: list[FlowNodeSchema] edges: list[FlowEdgeSchema] @@ -48,6 +73,8 @@ class FlowUpdate(BaseModel): class FlowCreate(BaseModel): + model_config = _FLOW_CONFIG + id: str name: str nodes: list[FlowNodeSchema] = [] diff --git a/api/v2/schemas/identities.py b/api/v2/schemas/identities.py index 57e2138b..857109cd 100644 --- a/api/v2/schemas/identities.py +++ b/api/v2/schemas/identities.py @@ -53,6 +53,7 @@ class ChangePasswordRequest(BaseModel): class LoginResponse(BaseModel): identity_id: str + identity: IdentityResponse class SendLoginCodeRequest(BaseModel): @@ -115,3 +116,17 @@ class LinkEmailConfirmRequest(BaseModel): class IdentityAttachTelegram(BaseModel): tg_id: int = Field(...) + + +class IdentitySessionItem(BaseModel): + id: str + device_label: str | None = None + ip: str | None = None + created_at: datetime + last_seen_at: datetime + expires_at: datetime | None = None + is_current: bool = False + + +class IdentitySessionsResponse(BaseModel): + sessions: list[IdentitySessionItem] diff --git a/api/v2/schemas/web.py b/api/v2/schemas/web.py index fc24f01c..22d0bbb5 100644 --- a/api/v2/schemas/web.py +++ b/api/v2/schemas/web.py @@ -31,12 +31,29 @@ class WebTheme(BaseModel): tokens: dict[str, Any] +class WebPageThemeResponse(BaseModel): + slug: str + variant_key: str = "default" + tokens: dict[str, Any] = Field(default_factory=dict) + + +class WebPageThemeUpdate(BaseModel): + tokens: dict[str, Any] + + class WebPageVariantSummary(BaseModel): key: str = Field(..., max_length=64) name: str = Field(..., max_length=255) is_active: bool = False +class WebPageSaveResponse(BaseModel): + slug: str + variant_key: str = "default" + active_variant_key: str = "default" + variants: list[WebPageVariantSummary] = Field(default_factory=list) + + class WebPageResponse(BaseModel): slug: str blocks: list[WebBlockResponse] diff --git a/database/__init__.py b/database/__init__.py index a310f320..7b919eae 100644 --- a/database/__init__.py +++ b/database/__init__.py @@ -1,4 +1,5 @@ from . import identities +from . import identity_sessions from .audit import * from .bans import * from .coupons import * diff --git a/database/identities.py b/database/identities.py index 3ed674a5..2df61c66 100644 --- a/database/identities.py +++ b/database/identities.py @@ -10,10 +10,30 @@ from sqlalchemy.ext.asyncio import AsyncSession from config import API_TOKEN_TTL_DAYS from core.executor import run_cpu, run_io +from database import identity_sessions as _idsess from database.access.tg_mirror import refresh_tg_mirrors_for_user from database.models import Admin, Identity, User +def _request_meta(request) -> tuple[str | None, str | None]: + if request is None: + return None, None + try: + ua = request.headers.get("user-agent") + except Exception: + ua = None + ip = None + try: + xff = request.headers.get("x-forwarded-for") + if xff: + ip = xff.split(",")[0].strip() + elif request.client and request.client.host: + ip = request.client.host + except Exception: + ip = None + return ua, ip + + _BCRYPT_MAX_PASSWORD_BYTES = 72 _BCRYPT_ROUNDS = 12 @@ -263,59 +283,68 @@ async def get_identity_by_token_hash(session: AsyncSession, token_hash: str) -> return result.scalar_one_or_none() -async def issue_token_for_identity(session: AsyncSession, identity: Identity) -> str: - """Генерирует токен, сохраняет хеш и token_issued_at в identity, возвращает токен (показать один раз).""" +async def issue_token_for_identity( + session: AsyncSession, + identity: Identity, + *, + request=None, +) -> str: + """Генерирует токен и создаёт новую сессию в identity_sessions (не затирая другие устройства).""" token = generate_token() - identity.api_token_hash = await run_io(hash_token, token) - identity.token_issued_at = datetime.utcnow() - await session.flush() + token_hash = await run_io(hash_token, token) + user_agent, ip = _request_meta(request) + await _idsess.create_identity_session( + session, + identity=identity, + token_hash=token_hash, + user_agent=user_agent, + ip=ip, + ) return token -def _is_token_expired(identity: Identity) -> bool: - """Проверяет, истёк ли срок действия токена (если задан API_TOKEN_TTL_DAYS).""" - if API_TOKEN_TTL_DAYS is None or identity.token_issued_at is None: - return False - expiry = identity.token_issued_at + timedelta(days=API_TOKEN_TTL_DAYS) - return datetime.utcnow() >= expiry - - async def create_identity_with_token( session: AsyncSession, email: str | None = None, password: str | None = None, tg_id: int | None = None, + *, + request=None, ) -> tuple[Identity, str]: """Создаёт идентичность и выдаёт API-токен. При регистрации по почте передать email и password.""" identity = await create_identity(session, email=email, tg_id=tg_id) if password: identity.password_hash = await run_cpu(hash_password, password) await session.refresh(identity) - token = await issue_token_for_identity(session, identity) + token = await issue_token_for_identity(session, identity, request=request) return identity, token async def verify_identity_token(session: AsyncSession, identity_id: str, token: str) -> Identity | None: - """Проверяет пару identity_id + token и срок действия токена; возвращает Identity или None.""" - identity = await get_identity_by_id(session, identity_id) - if not identity or not identity.api_token_hash: - return None + """Проверяет пару identity_id + token через identity_sessions; возвращает Identity или None.""" token_hash = await run_io(hash_token, token) - if token_hash != identity.api_token_hash: + sess = await _idsess.get_session_by_token_hash(session, token_hash) + if sess is None or sess.identity_id != identity_id: return None - if _is_token_expired(identity): + if sess.expires_at is not None and sess.expires_at <= datetime.utcnow(): return None - return identity + return await get_identity_by_id(session, identity_id) -async def login_by_email(session: AsyncSession, email: str, password: str) -> tuple[Identity, str] | None: +async def login_by_email( + session: AsyncSession, + email: str, + password: str, + *, + request=None, +) -> tuple[Identity, str] | None: """Вход по email и паролю: проверяет пароль, выдаёт новый токен; возвращает (identity, token) или None.""" identity = await get_identity_by_email(session, email) if not identity: return None if not await run_cpu(check_password, password, identity.password_hash): return None - token = await issue_token_for_identity(session, identity) + token = await issue_token_for_identity(session, identity, request=request) return identity, token diff --git a/database/identity_sessions.py b/database/identity_sessions.py new file mode 100644 index 00000000..63673d13 --- /dev/null +++ b/database/identity_sessions.py @@ -0,0 +1,149 @@ +from datetime import datetime, timedelta + +from sqlalchemy import delete, select, update +from sqlalchemy.ext.asyncio import AsyncSession + +from config import API_TOKEN_TTL_DAYS +from database.models import Identity, IdentitySession + + +_LAST_SEEN_TOUCH_SECONDS = 60 + + +def _device_label_from_user_agent(user_agent: str | None) -> str: + if not user_agent: + return "Unknown device" + ua = user_agent.lower() + browser = "Browser" + for needle, label in ( + ("edg/", "Edge"), + ("opr/", "Opera"), + ("yabrowser", "Yandex"), + ("firefox", "Firefox"), + ("chrome", "Chrome"), + ("safari", "Safari"), + ): + if needle in ua: + browser = label + break + platform = "Desktop" + if "iphone" in ua or "ipad" in ua or "ios" in ua: + platform = "iOS" + elif "android" in ua: + platform = "Android" + elif "macintosh" in ua or "mac os" in ua: + platform = "macOS" + elif "windows" in ua: + platform = "Windows" + elif "linux" in ua: + platform = "Linux" + return f"{browser} · {platform}"[:128] + + +async def create_identity_session( + session: AsyncSession, + *, + identity: Identity, + token_hash: str, + user_agent: str | None = None, + ip: str | None = None, + device_label: str | None = None, +) -> IdentitySession: + """Создаёт новую сессию. TTL — из API_TOKEN_TTL_DAYS.""" + now = datetime.utcnow() + expires_at = now + timedelta(days=API_TOKEN_TTL_DAYS) if API_TOKEN_TTL_DAYS else None + label = device_label or _device_label_from_user_agent(user_agent) + obj = IdentitySession( + identity_id=identity.id, + token_hash=token_hash, + device_label=label, + user_agent=user_agent, + ip=ip, + created_at=now, + last_seen_at=now, + expires_at=expires_at, + ) + session.add(obj) + await session.flush() + return obj + + +async def get_session_by_token_hash( + session: AsyncSession, token_hash: str +) -> IdentitySession | None: + result = await session.execute( + select(IdentitySession).where(IdentitySession.token_hash == token_hash) + ) + return result.scalar_one_or_none() + + +async def list_sessions_for_identity( + session: AsyncSession, identity_id: str +) -> list[IdentitySession]: + now = datetime.utcnow() + result = await session.execute( + select(IdentitySession) + .where(IdentitySession.identity_id == identity_id) + .where( + (IdentitySession.expires_at.is_(None)) | (IdentitySession.expires_at > now) + ) + .order_by(IdentitySession.last_seen_at.desc()) + ) + return list(result.scalars().all()) + + +async def delete_session_by_id( + session: AsyncSession, *, session_id: str, identity_id: str +) -> bool: + """Удаляет сессию по id при условии, что она принадлежит identity_id.""" + result = await session.execute( + delete(IdentitySession) + .where(IdentitySession.id == session_id) + .where(IdentitySession.identity_id == identity_id) + ) + return (result.rowcount or 0) > 0 + + +async def delete_session_by_token_hash(session: AsyncSession, token_hash: str) -> bool: + result = await session.execute( + delete(IdentitySession).where(IdentitySession.token_hash == token_hash) + ) + return (result.rowcount or 0) > 0 + + +async def delete_other_sessions( + session: AsyncSession, *, identity_id: str, keep_token_hash: str +) -> int: + """Удаляет все сессии identity кроме той, которая соответствует keep_token_hash.""" + result = await session.execute( + delete(IdentitySession) + .where(IdentitySession.identity_id == identity_id) + .where(IdentitySession.token_hash != keep_token_hash) + ) + return int(result.rowcount or 0) + + +async def touch_session_last_seen( + session: AsyncSession, sess: IdentitySession +) -> None: + """Обновляет last_seen_at, но только если прошло ≥60 секунд — чтобы не писать в БД на каждый запрос.""" + now = datetime.utcnow() + if (now - sess.last_seen_at).total_seconds() < _LAST_SEEN_TOUCH_SECONDS: + return + await session.execute( + update(IdentitySession) + .where(IdentitySession.id == sess.id) + .values(last_seen_at=now) + ) + sess.last_seen_at = now + + +async def cleanup_expired_sessions(session: AsyncSession) -> int: + """Удаляет все сессии с expires_at в прошлом. Возвращает количество удалённых.""" + now = datetime.utcnow() + result = await session.execute( + delete(IdentitySession) + .where(IdentitySession.expires_at.is_not(None)) + .where(IdentitySession.expires_at <= now) + ) + return int(result.rowcount or 0) diff --git a/database/migrations/schema_upgrade.py b/database/migrations/schema_upgrade.py index f372a802..e780b654 100644 --- a/database/migrations/schema_upgrade.py +++ b/database/migrations/schema_upgrade.py @@ -1225,6 +1225,58 @@ async def _migration_v23_add_identity_onboarding_stage(conn: AsyncConnection) -> await _exec_ignore(conn, "ALTER TABLE identities ADD COLUMN onboarding_stage VARCHAR(32)") +async def _migration_v24_add_identity_sessions(conn: AsyncConnection) -> None: + logger.info("[schema_upgrade] v24: таблица identity_sessions + перенос существующих токенов") + if not await _table_exists(conn, "identities"): + return + if not await _table_exists(conn, "identity_sessions"): + await _exec_ignore( + conn, + """ + CREATE TABLE identity_sessions ( + id VARCHAR(36) PRIMARY KEY, + identity_id VARCHAR(36) NOT NULL REFERENCES identities(id) ON DELETE CASCADE, + token_hash VARCHAR(64) NOT NULL UNIQUE, + device_label VARCHAR(128), + user_agent TEXT, + ip VARCHAR(64), + created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + last_seen_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP, + expires_at TIMESTAMP + ) + """, + ) + await _exec_ignore( + conn, + "CREATE INDEX IF NOT EXISTS ix_identity_sessions_identity_id ON identity_sessions(identity_id)", + ) + await _exec_ignore( + conn, + "CREATE INDEX IF NOT EXISTS ix_identity_sessions_identity_last_seen " + "ON identity_sessions(identity_id, last_seen_at)", + ) + if await _column_exists(conn, "identities", "api_token_hash"): + await _exec_ignore(conn, "CREATE EXTENSION IF NOT EXISTS pgcrypto") + await _exec_ignore( + conn, + """ + INSERT INTO identity_sessions ( + id, identity_id, token_hash, device_label, created_at, last_seen_at + ) + SELECT + gen_random_uuid()::text, + id, + api_token_hash, + 'legacy', + COALESCE(token_issued_at, CURRENT_TIMESTAMP), + COALESCE(token_issued_at, CURRENT_TIMESTAMP) + FROM identities + WHERE api_token_hash IS NOT NULL + AND api_token_hash NOT IN (SELECT token_hash FROM identity_sessions) + """, + ) + + _MIGRATIONS = [ (1, "Добавление users.id", _migration_v1_add_users_id), (2, "Добавление user_id колонок", _migration_v2_add_user_id_columns), @@ -1249,6 +1301,7 @@ _MIGRATIONS = [ (21, "identities.yandex_sub", _migration_v21_add_identity_yandex_sub), (22, "identities.onboarding_completed_at", _migration_v22_add_identity_onboarding_completed_at), (23, "identities.onboarding_stage", _migration_v23_add_identity_onboarding_stage), + (24, "таблица identity_sessions (мультидевайс)", _migration_v24_add_identity_sessions), ] diff --git a/database/models/__init__.py b/database/models/__init__.py index c3f5a316..d997f08f 100644 --- a/database/models/__init__.py +++ b/database/models/__init__.py @@ -4,6 +4,7 @@ from .audit import AuditEvent from .coupons import Coupon, CouponUsage from .gifts import Gift, GiftUsage from .identity import Identity +from .identity_session import IdentitySession from .keys import Key from .notifications import Notification, ScheduledBroadcast from .payments import Payment @@ -30,6 +31,7 @@ __all__ = [ "Base", "DictLikeMixin", "Identity", + "IdentitySession", "User", "ManualBan", "TemporaryData", diff --git a/database/models/identity_session.py b/database/models/identity_session.py new file mode 100644 index 00000000..92687146 --- /dev/null +++ b/database/models/identity_session.py @@ -0,0 +1,39 @@ +import uuid + +from datetime import datetime + +from sqlalchemy import ( + Column, + DateTime, + ForeignKey, + Index, + String, + Text, +) + +from ._base import Base, DictLikeMixin + + +class IdentitySession(DictLikeMixin, Base): + """Активная сессия identity. Один identity может иметь несколько сессий (devices).""" + + __tablename__ = "identity_sessions" + + id = Column(String(36), primary_key=True, default=lambda: str(uuid.uuid4())) + identity_id = Column( + String(36), + ForeignKey("identities.id", ondelete="CASCADE"), + nullable=False, + index=True, + ) + token_hash = Column(String(64), nullable=False, unique=True) + device_label = Column(String(128), nullable=True) + user_agent = Column(Text, nullable=True) + ip = Column(String(64), nullable=True) + created_at = Column(DateTime, default=datetime.utcnow, nullable=False) + last_seen_at = Column(DateTime, default=datetime.utcnow, nullable=False) + expires_at = Column(DateTime, nullable=True) + + __table_args__ = ( + Index("ix_identity_sessions_identity_last_seen", "identity_id", "last_seen_at"), + ) diff --git a/utils/versioning.py b/utils/versioning.py index ce5fe2ef..12d4a9dd 100644 --- a/utils/versioning.py +++ b/utils/versioning.py @@ -92,4 +92,4 @@ def get_git_commit_number() -> str: def get_version() -> str: - return f"v.6-b1804261111 {get_git_commit_number()}" + return f"v.6-b1904261111 {get_git_commit_number()}"