auth sessions + flow schema fixes
This commit is contained in:
+15
-3
@@ -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
|
||||
|
||||
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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] = []
|
||||
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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}
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
+52
-1
@@ -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)
|
||||
|
||||
|
||||
|
||||
+29
-2
@@ -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] = []
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
from . import identities
|
||||
from . import identity_sessions
|
||||
from .audit import *
|
||||
from .bans import *
|
||||
from .coupons import *
|
||||
|
||||
+52
-23
@@ -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
|
||||
|
||||
|
||||
|
||||
@@ -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)
|
||||
@@ -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),
|
||||
]
|
||||
|
||||
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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"),
|
||||
)
|
||||
+1
-1
@@ -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()}"
|
||||
|
||||
Reference in New Issue
Block a user