Реализовано без authlib, на голом httpx (AsyncClient — синхронный httpx
блокировал бы event loop API на время внешнего запроса), по образцу двух
провайдеров:
- Миграция 004: hashed_password → nullable (OAuth-юзеры без пароля),
oauth_provider/oauth_id + уникальный индекс на пару.
- app/core/security.py: verify_password защищён от hashed=None (иначе TypeError
при попытке OAuth-юзера войти по паролю — нашёл при ревью, не баг-репорт).
- app/core/oauth.py: get_authorize_url()/exchange_code() — единый интерфейс для
google/yandex. Redirect URI: <APP_URL>/api/auth/<provider>/callback.
- app/api/auth.py: GET /auth/{provider}/login (редирект на согласие, state в
httponly-cookie от CSRF) и /callback (обмен code, find-or-create юзера по
oauth_id → по email для привязки существующего аккаунта → новый без пароля,
is_verified=email_verified от провайдера). Токен фронту — через URL-фрагмент
#token=..., не query (не уходит в логи/Referer).
- Фронтенд: OAuthButtons (Login/Register), страница /oauth/callback (читает
фрагмент → GET /auth/me → setAuth → редирект в кабинет).
- 6 юнит-тестов чистой логики сборки ссылок (app/core/oauth.py) — первый тест-
контур для api/ в этой сессии (pytest.ini/conftest/requirements-test по
образцу остальных сервисов), добавлен в общий run_tests.sh + mypy-гейт.
GOOGLE_CLIENT_ID/SECRET уже в .env (юзер создал OAuth-клиент), YANDEX_* пусты —
эндпоинты в этом случае отвечают 503, не падают. .env.example документирует обе
пары. Тестов всего: 118 (было 112).
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
361 lines
14 KiB
Python
361 lines
14 KiB
Python
"""Роутер аутентификации: регистрация, вход, верификация email, OAuth."""
|
||
|
||
import logging
|
||
import secrets
|
||
|
||
from fastapi import APIRouter, Depends, HTTPException, Request, status
|
||
from fastapi.responses import RedirectResponse
|
||
from sqlalchemy import select
|
||
from sqlalchemy.ext.asyncio import AsyncSession
|
||
|
||
from app.config import settings
|
||
from app.core.celery_app import celery_app
|
||
from app.core.oauth import OAuthNotConfigured, OAuthUserInfo, exchange_code, get_authorize_url
|
||
from app.core.security import (
|
||
create_access_token,
|
||
get_current_user,
|
||
hash_password,
|
||
invalidate_user_cache,
|
||
verify_password,
|
||
)
|
||
from app.database import get_db
|
||
from app.models.user import User
|
||
from app.schemas.auth import (
|
||
ChangePasswordRequest,
|
||
LoginRequest,
|
||
RegisterRequest,
|
||
TokenResponse,
|
||
UpdateProfileRequest,
|
||
UserInToken,
|
||
)
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
router = APIRouter(prefix="/auth", tags=["auth"])
|
||
|
||
_OAUTH_PROVIDERS = {"google", "yandex"}
|
||
|
||
|
||
@router.post("/register", response_model=TokenResponse, status_code=status.HTTP_201_CREATED)
|
||
async def register(data: RegisterRequest, db: AsyncSession = Depends(get_db)) -> TokenResponse:
|
||
"""
|
||
Регистрация нового пользователя.
|
||
|
||
Создаёт аккаунт и отправляет письмо с подтверждением email.
|
||
"""
|
||
# Проверить уникальность email
|
||
existing = await db.execute(select(User).where(User.email == data.email.lower()))
|
||
if existing.scalar_one_or_none():
|
||
raise HTTPException(
|
||
status_code=status.HTTP_409_CONFLICT,
|
||
detail="Пользователь с таким email уже существует",
|
||
)
|
||
|
||
# Создать токен верификации
|
||
verification_token = secrets.token_urlsafe(32)
|
||
|
||
user = User(
|
||
email=data.email.lower(),
|
||
hashed_password=hash_password(data.password),
|
||
name=data.name,
|
||
verification_token=verification_token,
|
||
is_verified=False,
|
||
plan="free",
|
||
)
|
||
db.add(user)
|
||
await db.flush() # Получить ID без коммита
|
||
await db.commit()
|
||
await db.refresh(user)
|
||
|
||
# Диспатч email верификации через воркер
|
||
try:
|
||
celery_app.send_task(
|
||
"notify.send_verification",
|
||
args=[user.email, user.name, verification_token],
|
||
queue="queue.notify",
|
||
)
|
||
except Exception as e:
|
||
logger.warning(f"Не удалось поставить задачу верификации email: {e}")
|
||
|
||
access_token = create_access_token({"sub": str(user.id)})
|
||
return TokenResponse(
|
||
access_token=access_token,
|
||
token_type="bearer",
|
||
user=UserInToken.model_validate(user),
|
||
)
|
||
|
||
|
||
@router.post("/login", response_model=TokenResponse)
|
||
async def login(data: LoginRequest, db: AsyncSession = Depends(get_db)) -> TokenResponse:
|
||
"""Вход в систему. Возвращает JWT токен."""
|
||
result = await db.execute(select(User).where(User.email == data.email.lower()))
|
||
user = result.scalar_one_or_none()
|
||
|
||
if user is None or not verify_password(data.password, user.hashed_password):
|
||
raise HTTPException(
|
||
status_code=status.HTTP_401_UNAUTHORIZED,
|
||
detail="Неверный email или пароль",
|
||
)
|
||
|
||
access_token = create_access_token({"sub": str(user.id)})
|
||
return TokenResponse(
|
||
access_token=access_token,
|
||
token_type="bearer",
|
||
user=UserInToken.model_validate(user),
|
||
)
|
||
|
||
|
||
@router.get("/me", response_model=UserInToken)
|
||
async def get_me(current_user: User = Depends(get_current_user)) -> UserInToken:
|
||
"""Получить информацию о текущем пользователе."""
|
||
return UserInToken.model_validate(current_user)
|
||
|
||
|
||
@router.patch("/me", response_model=UserInToken)
|
||
async def update_profile(
|
||
data: UpdateProfileRequest,
|
||
current_user: User = Depends(get_current_user),
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> UserInToken:
|
||
"""
|
||
Изменить персональные данные (имя и/или email).
|
||
|
||
Смена email сбрасывает подтверждение и отправляет новое письмо верификации.
|
||
"""
|
||
# Загружаем "живого" пользователя из БД: объект из Redis-кэша не привязан
|
||
# к сессии и не содержит части полей.
|
||
result = await db.execute(select(User).where(User.id == current_user.id))
|
||
user = result.scalar_one_or_none()
|
||
if user is None:
|
||
raise HTTPException(
|
||
status_code=status.HTTP_404_NOT_FOUND, detail="Пользователь не найден"
|
||
)
|
||
|
||
email_changed = False
|
||
|
||
if data.name is not None:
|
||
user.name = data.name
|
||
|
||
if data.email is not None:
|
||
new_email = data.email.lower()
|
||
if new_email != user.email:
|
||
existing = await db.execute(
|
||
select(User).where(User.email == new_email, User.id != user.id)
|
||
)
|
||
if existing.scalar_one_or_none():
|
||
raise HTTPException(
|
||
status_code=status.HTTP_409_CONFLICT,
|
||
detail="Этот email уже используется другим аккаунтом",
|
||
)
|
||
user.email = new_email
|
||
user.is_verified = False
|
||
user.verification_token = secrets.token_urlsafe(32)
|
||
email_changed = True
|
||
|
||
await db.commit()
|
||
await db.refresh(user)
|
||
await invalidate_user_cache(user.id)
|
||
|
||
# При смене email — отправить новое письмо подтверждения на новый адрес
|
||
if email_changed:
|
||
try:
|
||
celery_app.send_task(
|
||
"notify.send_verification",
|
||
args=[user.email, user.name, user.verification_token],
|
||
queue="queue.notify",
|
||
)
|
||
except Exception as e:
|
||
logger.warning(f"Не удалось отправить письмо верификации при смене email: {e}")
|
||
|
||
return UserInToken.model_validate(user)
|
||
|
||
|
||
@router.post("/change-password", status_code=status.HTTP_204_NO_CONTENT)
|
||
async def change_password(
|
||
data: ChangePasswordRequest,
|
||
current_user: User = Depends(get_current_user),
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> None:
|
||
"""Сменить пароль. Требует текущий пароль для подтверждения."""
|
||
result = await db.execute(select(User).where(User.id == current_user.id))
|
||
user = result.scalar_one_or_none()
|
||
if user is None:
|
||
raise HTTPException(
|
||
status_code=status.HTTP_404_NOT_FOUND, detail="Пользователь не найден"
|
||
)
|
||
|
||
if not verify_password(data.current_password, user.hashed_password):
|
||
raise HTTPException(
|
||
status_code=status.HTTP_400_BAD_REQUEST,
|
||
detail="Текущий пароль указан неверно",
|
||
)
|
||
|
||
if verify_password(data.new_password, user.hashed_password):
|
||
raise HTTPException(
|
||
status_code=status.HTTP_400_BAD_REQUEST,
|
||
detail="Новый пароль совпадает с текущим",
|
||
)
|
||
|
||
user.hashed_password = hash_password(data.new_password)
|
||
await db.commit()
|
||
await invalidate_user_cache(user.id)
|
||
|
||
|
||
@router.post("/resend-verification", status_code=status.HTTP_204_NO_CONTENT)
|
||
async def resend_verification(
|
||
current_user: User = Depends(get_current_user),
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> None:
|
||
"""Повторно отправить письмо с подтверждением email."""
|
||
if current_user.is_verified:
|
||
raise HTTPException(
|
||
status_code=status.HTTP_400_BAD_REQUEST,
|
||
detail="Email уже подтверждён",
|
||
)
|
||
|
||
if not current_user.verification_token:
|
||
current_user.verification_token = secrets.token_urlsafe(32)
|
||
await db.commit()
|
||
await db.refresh(current_user)
|
||
|
||
try:
|
||
celery_app.send_task(
|
||
"notify.send_verification",
|
||
args=[current_user.email, current_user.name, current_user.verification_token],
|
||
queue="queue.notify",
|
||
)
|
||
except Exception as e:
|
||
logger.warning(f"Не удалось поставить задачу повторной верификации: {e}")
|
||
raise HTTPException(
|
||
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
|
||
detail="Не удалось отправить письмо, попробуйте позже",
|
||
) from e
|
||
|
||
|
||
@router.post("/verify-email/{token}", status_code=status.HTTP_200_OK)
|
||
async def verify_email(token: str, db: AsyncSession = Depends(get_db)) -> dict:
|
||
"""Подтвердить email по токену из письма."""
|
||
result = await db.execute(
|
||
select(User).where(User.verification_token == token)
|
||
)
|
||
user = result.scalar_one_or_none()
|
||
|
||
if user is None:
|
||
raise HTTPException(
|
||
status_code=status.HTTP_404_NOT_FOUND,
|
||
detail="Неверный или просроченный токен верификации",
|
||
)
|
||
|
||
user.is_verified = True
|
||
user.verification_token = None
|
||
await db.commit()
|
||
await invalidate_user_cache(user.id)
|
||
|
||
return {"message": "Email успешно подтверждён"}
|
||
|
||
|
||
async def _get_or_create_oauth_user(
|
||
db: AsyncSession, provider: str, info: OAuthUserInfo
|
||
) -> User:
|
||
"""Найти пользователя по (provider, oauth_id), иначе по email (привязка
|
||
существующего аккаунта), иначе создать нового без пароля."""
|
||
result = await db.execute(
|
||
select(User).where(User.oauth_provider == provider, User.oauth_id == info.provider_id)
|
||
)
|
||
user = result.scalar_one_or_none()
|
||
if user:
|
||
return user
|
||
|
||
email = info.email.lower()
|
||
result = await db.execute(select(User).where(User.email == email))
|
||
user = result.scalar_one_or_none()
|
||
if user:
|
||
if not user.oauth_provider:
|
||
user.oauth_provider = provider
|
||
user.oauth_id = info.provider_id
|
||
if info.email_verified:
|
||
user.is_verified = True
|
||
await db.commit()
|
||
await db.refresh(user)
|
||
return user
|
||
|
||
user = User(
|
||
email=email,
|
||
hashed_password=None,
|
||
name=info.name,
|
||
is_verified=info.email_verified,
|
||
plan="free",
|
||
oauth_provider=provider,
|
||
oauth_id=info.provider_id,
|
||
)
|
||
db.add(user)
|
||
await db.commit()
|
||
await db.refresh(user)
|
||
return user
|
||
|
||
|
||
@router.get("/{provider}/login", include_in_schema=False)
|
||
async def oauth_login(provider: str) -> RedirectResponse:
|
||
"""Редирект на экран согласия Google/Яндекс."""
|
||
if provider not in _OAUTH_PROVIDERS:
|
||
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Неизвестный провайдер")
|
||
|
||
state = secrets.token_urlsafe(24)
|
||
try:
|
||
url = get_authorize_url(provider, state)
|
||
except OAuthNotConfigured as e:
|
||
raise HTTPException(
|
||
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
|
||
detail=f"Вход через {provider} временно недоступен",
|
||
) from e
|
||
|
||
response = RedirectResponse(url)
|
||
response.set_cookie(
|
||
f"oauth_state_{provider}", state,
|
||
httponly=True, secure=True, samesite="lax", max_age=600,
|
||
)
|
||
return response
|
||
|
||
|
||
@router.get("/{provider}/callback", include_in_schema=False)
|
||
async def oauth_callback(
|
||
provider: str,
|
||
request: Request,
|
||
code: str | None = None,
|
||
state: str | None = None,
|
||
db: AsyncSession = Depends(get_db),
|
||
) -> RedirectResponse:
|
||
"""Колбэк провайдера: обменять code на профиль, найти/создать юзера, выдать JWT.
|
||
|
||
Токен передаётся фронту через URL-фрагмент (#token=...) — он не уходит на
|
||
сервер при последующих запросах и не попадает в логи/Referer, в отличие от
|
||
query-параметра. Фронт (страница /oauth/callback) читает его и вызывает
|
||
setAuth, как после обычного /login.
|
||
"""
|
||
if provider not in _OAUTH_PROVIDERS:
|
||
raise HTTPException(status_code=status.HTTP_404_NOT_FOUND, detail="Неизвестный провайдер")
|
||
|
||
expected_state = request.cookies.get(f"oauth_state_{provider}")
|
||
if not code or not state or not expected_state or state != expected_state:
|
||
raise HTTPException(status_code=status.HTTP_400_BAD_REQUEST, detail="Невалидный OAuth-колбэк")
|
||
|
||
try:
|
||
info = await exchange_code(provider, code)
|
||
except OAuthNotConfigured as e:
|
||
raise HTTPException(
|
||
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
|
||
detail=f"Вход через {provider} временно недоступен",
|
||
) from e
|
||
except Exception as e:
|
||
logger.warning(f"OAuth {provider}: обмен кода не удался: {e}")
|
||
raise HTTPException(
|
||
status_code=status.HTTP_400_BAD_REQUEST, detail="Не удалось войти через провайдера"
|
||
) from e
|
||
|
||
user = await _get_or_create_oauth_user(db, provider, info)
|
||
access_token = create_access_token({"sub": str(user.id)})
|
||
|
||
response = RedirectResponse(f"{settings.APP_URL}/oauth/callback#token={access_token}")
|
||
response.delete_cookie(f"oauth_state_{provider}")
|
||
return response
|