Files
server-core/app/api/dependencies.py
T
Pine dbc29eee59 feat(membership): 成员制归属重构与叠加制身份(多对多成员表+转园两级审核)
- park_members/company_members 多对多成员表(0028),users.affiliation 为派生缓存由 sync_user_affiliation 刷新
- 转园两级审核状态机(0029):pending_from→pending_to→approved,支持用户/企业整体转园与运营端兜底
- 叠加制身份(0030):身份恒为运营/载体/OPC,企业管理员为能力叠加非身份,去除 enterprise/carrier 直写角色
- rbac 新增 memberships/park-transfers/企业管理员工作台 /admin/ent/* 接口

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-31 20:36:26 +08:00

141 lines
5.3 KiB
Python

# -*- coding: utf-8 -*-
"""接口层依赖:数据库句柄与当前登录用户(JWT + RBAC,全异步)。
``Database`` 实例由 ``main.py`` 在启动时创建并挂在 ``app.state.db`` 上,
路由通过 ``Depends(get_db)`` 取用;测试时可替换为临时目录实例。
"""
from __future__ import annotations
from fastapi import Depends, HTTPException, Request
from .. import config
from ..infrastructure.repositories import Database
from ..jwt import decode_access_token
BEARER_PREFIX = "Bearer "
def extract_bearer_token(request: Request) -> str:
"""从 Authorization 头提取 Bearer token,没有则返回空串。"""
auth_header = request.headers.get("Authorization", "")
if auth_header.startswith(BEARER_PREFIX):
return auth_header[len(BEARER_PREFIX):].strip()
return ""
async def get_db():
"""每请求一个独立 Database/AsyncSession(绑定共享引擎)。
之前直接返回 app.state.db(单例 AsyncSession)被并发请求复用,会间歇报
``sqlalchemy.exc.InvalidRequestError: concurrent operations are not permitted``。
改为每请求各建一个 Database 并在请求结束关闭,根除该竞态。
"""
db = Database()
try:
yield db
finally:
await db.close()
async def _compute_capabilities(db: Database, user: dict) -> list[str]:
"""叠加制能力集合(身份不互斥):
- opc_member:所有账号的基础能力;
- operator:运营端角色;
- carrier:园区端角色 或 已绑定园区管理员(park_members.admin / operator_user_id);
- enterprise:已绑定企业管理员(company_members.is_admin)。
"""
from ..infrastructure.models import CompanyMember, ParkMember, ParkTenant
from sqlalchemy import select
role = user.get("role", "")
caps = ["opc_member"]
if role in ("operator", "operator_internal", "op_admin", "op_super_admin", "superadmin", "admin"):
caps.append("operator")
if role in ("carrier", "park", "park_staff", "carrier_staff"):
caps.append("carrier")
else:
bound = await db.session.scalar(select(ParkMember).where(
ParkMember.user_id == user.get("id", ""), ParkMember.park_id != "",
ParkMember.member_type == "admin", ParkMember.status == "active").limit(1))
if not bound:
bound = await db.session.scalar(select(ParkTenant).where(
ParkTenant.operator_user_id == user.get("id", "")).limit(1))
if bound:
caps.append("carrier")
ent = await db.session.scalar(select(CompanyMember).where(
CompanyMember.user_id == user.get("id", ""), CompanyMember.is_admin.is_(True),
CompanyMember.status == "active").limit(1))
if ent:
caps.append("enterprise")
return caps
async def _resolve_identity(user: dict, db: Database) -> dict:
"""叠加制能力模型:role 为基础角色,capabilities 为绑定叠加出的能力集合。"""
user["permissions"] = await db.roles.permissions_for(
user.get("role", "opc_member"), None,
)
user["scope_region_ids"] = await db.regions.visible_region_ids(user.get("region_id"))
user["scope_level"] = await db.regions.level(user.get("region_id"))
user["capabilities"] = await _compute_capabilities(db, user)
return user
async def get_current_user(
request: Request,
db: Database = Depends(get_db),
) -> dict:
"""校验调用方 JWT,返回当前用户记录(含 role/权限/数据范围)。"""
token = extract_bearer_token(request)
if not token:
raise HTTPException(status_code=401, detail="No token provided")
payload = decode_access_token(token)
if payload is None:
raise HTTPException(status_code=401, detail="Invalid or expired token")
user_id = payload.get("sub")
if not await db.tokens.session_valid(
payload.get("jti", ""),
user_id,
payload.get("ver", 0),
):
raise HTTPException(status_code=401, detail="Invalid or expired token")
user = await db.users.get_by_id(user_id)
if user is None or user.get("status") != "active":
raise HTTPException(status_code=401, detail="Invalid or expired token")
return await _resolve_identity(user, db)
async def optional_current_user(
request: Request,
db: Database = Depends(get_db),
) -> dict | None:
"""可选登录:有有效 JWT 则返回用户,否则返回 None(不报错)。"""
token = extract_bearer_token(request)
if not token:
return None
payload = decode_access_token(token)
if payload is None:
return None
if not await db.tokens.session_valid(
payload.get("jti", ""), payload.get("sub", ""), payload.get("ver", 0),
):
return None
user = await db.users.get_by_id(payload.get("sub"))
if user is None or user.get("status") != "active":
return None
return await _resolve_identity(user, db)
def require_port(user: dict = Depends(get_current_user)) -> dict:
"""按端口隔离的资源(如智能体)—— 单角色后不再按端口隔离,仅要求已登录。
账号唯一角色:智能体不再按端口隔离,统一归到该账号(port="")。为兼容
既有 AgentRepository 的 port 过滤,这里补一个恒定空串。
"""
user.setdefault("port", "")
return user