feat(auth): 新增微信同步登录接口 wechat-sync(小程序BFF用) #185
@@ -6,13 +6,14 @@ app.dependencies and authentication behavior lives in application use cases.
|
||||
"""
|
||||
|
||||
import logging
|
||||
import os
|
||||
from typing import Optional
|
||||
|
||||
import jwt
|
||||
from app.auth import AuthenticatedUser, blacklist_token, get_current_user
|
||||
from app.config import settings
|
||||
from app.dependencies import get_auth_email_service, get_auth_session_store, get_user_repository
|
||||
from fastapi import APIRouter, Depends, HTTPException, status
|
||||
from fastapi import APIRouter, Depends, Header, HTTPException, status
|
||||
from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer
|
||||
from pydantic import BaseModel, EmailStr
|
||||
|
||||
@@ -284,3 +285,99 @@ def _translate_auth_error(error: str | None) -> str:
|
||||
"Display name is required": "显示名称不能为空",
|
||||
}
|
||||
return translations.get(error or "", error or "注册失败")
|
||||
|
||||
|
||||
class WechatSyncRequest(BaseModel):
|
||||
openid: str
|
||||
unionid: Optional[str] = None
|
||||
nickname: Optional[str] = None
|
||||
avatar_url: Optional[str] = None
|
||||
source: str = "miniapp"
|
||||
|
||||
|
||||
class WechatSyncResponse(BaseModel):
|
||||
access_token: str
|
||||
token: str
|
||||
refresh_token: str
|
||||
user_id: str
|
||||
user: dict
|
||||
user_info: dict
|
||||
is_new_user: bool
|
||||
expires_in: int
|
||||
|
||||
|
||||
def _get_internal_api_keys() -> list[str]:
|
||||
"""获取内部 API Key 列表
|
||||
|
||||
优先级:
|
||||
1. INTERNAL_API_KEYS 环境变量
|
||||
2. /app/generated/internal_api_keys.txt 文件 (volume 持久化)
|
||||
"""
|
||||
env_keys = os.environ.get("INTERNAL_API_KEYS", "")
|
||||
if env_keys:
|
||||
return [k.strip() for k in env_keys.split(",") if k.strip()]
|
||||
|
||||
# 从持久化文件读取
|
||||
try:
|
||||
with open("/app/generated/internal_api_keys.txt", "r") as f:
|
||||
content = f.read().strip()
|
||||
if content:
|
||||
return [k.strip() for k in content.split(",") if k.strip()]
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
return []
|
||||
|
||||
|
||||
def _verify_internal_api_key(x_api_key: str | None = Header(None)) -> bool:
|
||||
"""验证内部 API Key
|
||||
|
||||
- 已配置时:必须匹配 INTERNAL_API_KEYS 中的 key
|
||||
- 未配置且非生产环境:放行(方便开发)
|
||||
- 未配置且生产环境:拒绝
|
||||
"""
|
||||
env = os.environ.get("APP_ENV", os.environ.get("ENV", "development")).lower()
|
||||
key_list = _get_internal_api_keys()
|
||||
|
||||
if not key_list:
|
||||
if env in ("production", "prod"):
|
||||
raise HTTPException(status_code=401, detail="内部接口未配置 API Key")
|
||||
return True
|
||||
|
||||
if x_api_key and x_api_key.strip() in key_list:
|
||||
return True
|
||||
|
||||
raise HTTPException(status_code=401, detail="无效的 API Key")
|
||||
|
||||
|
||||
@router.post("/wechat-sync", response_model=WechatSyncResponse, include_in_schema=False)
|
||||
async def wechat_sync(
|
||||
request: WechatSyncRequest,
|
||||
user_repository: UserRepository = Depends(get_user_repository),
|
||||
_: bool = Depends(_verify_internal_api_key),
|
||||
):
|
||||
"""
|
||||
微信同步登录/注册(系统级内部接口)
|
||||
|
||||
由 BFF 层通过 API Key 调用,不直接面向终端用户。
|
||||
根据 openid 查找或创建用户,返回 SaaS token。
|
||||
"""
|
||||
from packages.application.auth.wechat_sync_use_case import (
|
||||
WechatSyncRequest as UseCaseRequest,
|
||||
WechatSyncUseCase,
|
||||
)
|
||||
|
||||
use_case = WechatSyncUseCase(user_repository=user_repository)
|
||||
use_case_request = UseCaseRequest(
|
||||
openid=request.openid,
|
||||
unionid=request.unionid,
|
||||
nickname=request.nickname,
|
||||
avatar_url=request.avatar_url,
|
||||
source=request.source,
|
||||
)
|
||||
|
||||
response, error = use_case.execute(use_case_request)
|
||||
if error:
|
||||
raise HTTPException(status_code=400, detail=error)
|
||||
|
||||
return WechatSyncResponse(**response.to_dict())
|
||||
|
||||
@@ -20,6 +20,7 @@ from app.auth import AuthenticatedUser, get_current_user
|
||||
from app.dependencies import get_db_session
|
||||
from app.services import EditTemplateService
|
||||
from fastapi import APIRouter, Depends, HTTPException, Query, status
|
||||
from fastapi.responses import Response
|
||||
from pydantic import BaseModel, Field
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
@@ -263,7 +264,7 @@ def delete_template(
|
||||
template_id: str,
|
||||
db: Session = Depends(get_db_session),
|
||||
current_user: AuthenticatedUser = Depends(get_current_user),
|
||||
) -> None:
|
||||
) -> Response:
|
||||
"""删除模板(软删除 → 设为 inactive)"""
|
||||
_require_admin(current_user)
|
||||
svc = EditTemplateService(db)
|
||||
@@ -275,3 +276,4 @@ def delete_template(
|
||||
detail=str(exc),
|
||||
)
|
||||
logger.info("删除模板(软删除): id=%s by user=%s", template_id, current_user.user.id)
|
||||
return Response(status_code=204)
|
||||
|
||||
@@ -33,6 +33,8 @@ class SQLAlchemyUserRepository(UserRepository):
|
||||
model.max_projects = user.max_projects
|
||||
model.max_storage_gb = user.max_storage_gb
|
||||
model.is_admin = user.is_admin
|
||||
model.wechat_openid = user.wechat_openid
|
||||
model.wechat_unionid = user.wechat_unionid
|
||||
model.created_at = user.created_at
|
||||
|
||||
self.session.commit()
|
||||
@@ -49,6 +51,16 @@ class SQLAlchemyUserRepository(UserRepository):
|
||||
model = self.session.query(UserModel).filter(UserModel.username == username.strip()).first()
|
||||
return self._to_entity(model)
|
||||
|
||||
def find_by_wechat_openid(self, openid: str) -> User | None:
|
||||
model = self.session.query(UserModel).filter(UserModel.wechat_openid == openid.strip()).first()
|
||||
return self._to_entity(model)
|
||||
|
||||
def find_by_wechat_unionid(self, unionid: str) -> User | None:
|
||||
if not unionid or not unionid.strip():
|
||||
return None
|
||||
model = self.session.query(UserModel).filter(UserModel.wechat_unionid == unionid.strip()).first()
|
||||
return self._to_entity(model)
|
||||
|
||||
def find_by_verification_token(self, token: str) -> User | None:
|
||||
model = self.session.query(UserModel).filter(UserModel.email_verification_token == token).first()
|
||||
return self._to_entity(model)
|
||||
@@ -87,5 +99,7 @@ class SQLAlchemyUserRepository(UserRepository):
|
||||
max_projects=model.max_projects or 3,
|
||||
max_storage_gb=model.max_storage_gb or 10,
|
||||
is_admin=model.is_admin or False,
|
||||
wechat_openid=model.wechat_openid,
|
||||
wechat_unionid=model.wechat_unionid,
|
||||
created_at=model.created_at,
|
||||
)
|
||||
|
||||
@@ -0,0 +1,205 @@
|
||||
"""
|
||||
微信同步登录/注册 Use Case
|
||||
|
||||
供 BFF 层调用的系统级接口:
|
||||
- 根据 openid 查找用户,找到则登录返回 token
|
||||
- 没找到则创建新用户并返回 token
|
||||
- 支持 unionid 跨应用关联
|
||||
"""
|
||||
|
||||
import secrets
|
||||
from datetime import datetime, timedelta, timezone
|
||||
from typing import Optional
|
||||
from uuid import uuid4
|
||||
|
||||
import jwt as pyjwt
|
||||
|
||||
from packages.adapters.redis import get_session_store
|
||||
from packages.application.auth.jwt_service import jwt_service
|
||||
from packages.domain.entities import User
|
||||
|
||||
|
||||
class WechatSyncRequest:
|
||||
"""微信同步登录请求"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
openid: str,
|
||||
unionid: str = "",
|
||||
nickname: str = "",
|
||||
avatar_url: str = "",
|
||||
source: str = "miniapp",
|
||||
):
|
||||
self.openid = openid.strip()
|
||||
self.unionid = unionid.strip() if unionid else ""
|
||||
self.nickname = nickname or "微信用户"
|
||||
self.avatar_url = avatar_url or ""
|
||||
self.source = source
|
||||
|
||||
|
||||
class WechatSyncResponse:
|
||||
"""微信同步登录响应"""
|
||||
|
||||
def __init__(
|
||||
self,
|
||||
access_token: str,
|
||||
refresh_token: str,
|
||||
user_id: str,
|
||||
nickname: str,
|
||||
avatar_url: str,
|
||||
is_new_user: bool,
|
||||
expires_in: int,
|
||||
):
|
||||
self.access_token = access_token
|
||||
self.refresh_token = refresh_token
|
||||
self.user_id = user_id
|
||||
self.nickname = nickname
|
||||
self.avatar_url = avatar_url
|
||||
self.is_new_user = is_new_user
|
||||
self.expires_in = expires_in
|
||||
|
||||
def to_dict(self) -> dict:
|
||||
return {
|
||||
"access_token": self.access_token,
|
||||
"token": self.access_token, # 兼容 BFF 层用 token 字段读取
|
||||
"refresh_token": self.refresh_token,
|
||||
"user_id": self.user_id,
|
||||
"user": {
|
||||
"id": self.user_id,
|
||||
"nickname": self.nickname,
|
||||
"avatar_url": self.avatar_url,
|
||||
"display_name": self.nickname,
|
||||
},
|
||||
"user_info": {
|
||||
"id": self.user_id,
|
||||
"nickname": self.nickname,
|
||||
"avatar_url": self.avatar_url,
|
||||
"display_name": self.nickname,
|
||||
},
|
||||
"is_new_user": self.is_new_user,
|
||||
"expires_in": self.expires_in,
|
||||
}
|
||||
|
||||
|
||||
class WechatSyncUseCase:
|
||||
"""微信同步登录/注册用例
|
||||
|
||||
系统级接口,由 BFF 通过 API Key 调用。
|
||||
职责:根据 openid 查找或创建用户,返回 SaaS token。
|
||||
"""
|
||||
|
||||
def __init__(self, user_repository, session_store=None, jwt_secret_key: str | None = None):
|
||||
self.user_repository = user_repository
|
||||
self.session_store = session_store or get_session_store()
|
||||
self.jwt_secret_key = jwt_secret_key or jwt_service.config.SECRET_KEY
|
||||
|
||||
def execute(self, request: WechatSyncRequest) -> tuple[Optional[WechatSyncResponse], Optional[str]]:
|
||||
"""
|
||||
执行微信同步登录/注册
|
||||
|
||||
Returns:
|
||||
(响应对象, 错误信息) - 成功则错误信息为 None
|
||||
"""
|
||||
try:
|
||||
if not request.openid:
|
||||
return None, "openid is required"
|
||||
|
||||
is_new_user = False
|
||||
|
||||
# 1. 按 openid 查找用户
|
||||
user = self.user_repository.find_by_wechat_openid(request.openid)
|
||||
|
||||
# 2. 如果 openid 没找到,尝试 unionid
|
||||
if not user and request.unionid:
|
||||
user = self.user_repository.find_by_wechat_unionid(request.unionid)
|
||||
if user:
|
||||
# 找到用户但 openid 为空,绑定一下当前 openid
|
||||
user.wechat_openid = request.openid
|
||||
self.user_repository.save(user)
|
||||
|
||||
# 3. 都没找到则创建新用户
|
||||
if not user:
|
||||
user = self._create_wechat_user(request)
|
||||
is_new_user = True
|
||||
|
||||
# 4. 创建 session 并生成 token
|
||||
session_id = secrets.token_urlsafe(16)
|
||||
refresh_token = secrets.token_urlsafe(32)
|
||||
|
||||
now = datetime.now(timezone.utc)
|
||||
access_token_payload = {
|
||||
"sub": user.id,
|
||||
"sid": session_id,
|
||||
"type": "user_auth",
|
||||
"iat": now,
|
||||
"exp": now + timedelta(minutes=jwt_service.config.ACCESS_TOKEN_EXPIRE_MINUTES),
|
||||
}
|
||||
access_token = pyjwt.encode(
|
||||
access_token_payload,
|
||||
self.jwt_secret_key,
|
||||
algorithm=jwt_service.config.ALGORITHM,
|
||||
)
|
||||
|
||||
# 保存 session
|
||||
self.session_store.save_session(
|
||||
session_id=session_id,
|
||||
user_id=user.id,
|
||||
refresh_token=refresh_token,
|
||||
device_info=f"wechat_{request.source}",
|
||||
ip_address="bff_gateway",
|
||||
expires_in_seconds=30 * 24 * 3600, # 30 天
|
||||
)
|
||||
|
||||
# 更新最后登录信息
|
||||
user.last_login_at = now
|
||||
user.last_login_ip = "bff_gateway"
|
||||
self.user_repository.save(user)
|
||||
|
||||
response = WechatSyncResponse(
|
||||
access_token=access_token,
|
||||
refresh_token=refresh_token,
|
||||
user_id=user.id,
|
||||
nickname=user.display_name,
|
||||
avatar_url="", # SaaS 用户模型暂存头像,后续可扩展
|
||||
is_new_user=is_new_user,
|
||||
expires_in=jwt_service.config.ACCESS_TOKEN_EXPIRE_MINUTES * 60,
|
||||
)
|
||||
|
||||
return response, None
|
||||
|
||||
except Exception as e:
|
||||
return None, f"Internal error: {str(e)}"
|
||||
|
||||
def _create_wechat_user(self, request: WechatSyncRequest) -> User:
|
||||
"""创建微信用户"""
|
||||
user_id = uuid4().hex
|
||||
|
||||
# 生成唯一名和邮箱(微信用户无真实邮箱,用 openid 生成占位)
|
||||
safe_openid = request.openid.replace("-", "_")[:20]
|
||||
username = f"wx_{safe_openid}"
|
||||
email = f"{safe_openid}@wechat.local"
|
||||
|
||||
# 确保 username 唯一
|
||||
suffix = 0
|
||||
while self.user_repository.find_by_username(username):
|
||||
suffix += 1
|
||||
username = f"wx_{safe_openid}_{suffix}"
|
||||
|
||||
# 随机密码(微信用户不用密码登录)
|
||||
random_password = secrets.token_urlsafe(32)
|
||||
from packages.application.auth.password_hasher import password_hasher
|
||||
password_hash = password_hasher.hash_password(random_password)
|
||||
|
||||
user = User(
|
||||
id=user_id,
|
||||
email=email,
|
||||
username=username,
|
||||
display_name=request.nickname or "微信用户",
|
||||
password_hash=password_hash,
|
||||
email_verified=True, # 微信登录视为已验证
|
||||
wechat_openid=request.openid,
|
||||
wechat_unionid=request.unionid or None,
|
||||
)
|
||||
|
||||
self.user_repository.save(user)
|
||||
return user
|
||||
@@ -54,6 +54,11 @@ class User:
|
||||
used_storage_gb: float = 0.0
|
||||
# 管理员标识
|
||||
is_admin: bool = False
|
||||
|
||||
# 微信绑定
|
||||
wechat_openid: str | None = None
|
||||
wechat_unionid: str | None = None
|
||||
|
||||
created_at: datetime = field(default_factory=lambda: datetime.now(timezone.utc))
|
||||
|
||||
|
||||
|
||||
@@ -41,6 +41,16 @@ class UserRepository(ABC):
|
||||
"""根据密码重置令牌查找用户"""
|
||||
pass
|
||||
|
||||
@abstractmethod
|
||||
def find_by_wechat_openid(self, openid: str) -> Optional[User]:
|
||||
"""根据微信 openid 查找用户"""
|
||||
pass
|
||||
|
||||
@abstractmethod
|
||||
def find_by_wechat_unionid(self, unionid: str) -> Optional[User]:
|
||||
"""根据微信 unionid 查找用户"""
|
||||
pass
|
||||
|
||||
@abstractmethod
|
||||
def delete(self, user_id: str) -> bool:
|
||||
"""删除用户"""
|
||||
|
||||
Reference in New Issue
Block a user