Files
server-core/app/api/routers/rbac_admin.py
T

508 lines
21 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# -*- coding: utf-8 -*-
"""运营方管理端点(证明 RBAC 角色/权限/审计)。"""
from __future__ import annotations
from fastapi import APIRouter, Depends, HTTPException, Request
from pydantic import BaseModel, Field
from ..dependencies import get_db
from ..schemas.admin import SetUserRoleRequest, SetUserStatusRequest, UserClassificationRequest
from ..schemas.admin import UserCreateRequest, RolePermissionRequest, UpdateUserRequest
from ..schemas.admin import OpcCertReviewRequest, ParkAdmissionReviewRequest, ParkTransferReviewRequest
from ...rbac import require_permission, require_roles, write_audit
from ...infrastructure.repositories import Database
from ...services import compute_client
router = APIRouter(prefix="/admin", tags=["admin"])
async def _sync_compute(coro) -> None:
"""算力引擎用户同步为附加能力:失败不阻断平台用户操作(引擎未配置/不可达时静默)。"""
try:
await coro
except Exception: # noqa: BLE001
pass
@router.post("/users", summary="新增账号(按类型/角色)")
async def create_user(
req: UserCreateRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.assign_role")),
):
from ...services.user_admin_service import UserAdminService
user = await UserAdminService(db).create_user(req, actor)
# 平台用户创建 → 引擎同步建号 + 签发 PAT(算力引擎用户与平台一致)
await _sync_compute(compute_client.create_user(user["username"]))
await _sync_compute(compute_client.issue_pat(user["username"]))
await write_audit(db, action="user.create", resource="user", resource_id=user["id"],
detail=f"role={req.role}", user=actor, request=request)
return await db.users.to_profile(user)
def _port_for_role(role: str) -> str:
return {
"opc_member": "opc", "carrier": "carrier", "enterprise": "enterprise",
"provider": "provider", "government": "government", "operator": "operator",
"investor": "investor", "developer": "developer",
}.get(role, "opc")
@router.get("/users", summary="用户列表(运营方)")
async def list_users(
request: Request,
db: Database = Depends(get_db),
_role: dict = Depends(require_roles("operator")),
_perm: dict = Depends(require_permission("menu:admin_user_mgmt")),
):
return await db.users.list()
@router.post("/users/{user_id}/role", summary="分配角色")
async def set_user_role(
user_id: str,
req: SetUserRoleRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.assign_role")),
):
from ...services.user_admin_service import UserAdminService
if await db.users.get_by_id(user_id) is None:
raise HTTPException(status_code=404, detail="User not found")
updated = await UserAdminService(db).set_user_role(
actor, user_id, req.role, req.sub_role, req.org_id, req.region_id,
)
await write_audit(
db, action="role.assign", resource="user", resource_id=user_id,
detail=f"{actor['username']} -> role={req.role} sub={req.sub_role}",
user=actor, request=request,
)
return {"ok": True, "user": await db.users.to_profile(updated)}
@router.post("/users/{user_id}/classification", summary="设置用户认证状态/所属/账号类型")
async def set_user_classification(
user_id: str,
req: UserClassificationRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.manage")),
):
# 铁律:任何端口不得手动修改用户「所属园区」。affiliation/park_id/park_name 仅经
# 入驻申请 / 转园申请(审核通过)回填,此处剔除,只允许改 认证状态/账号类型。
fields = {k: v for k, v in req.model_dump(exclude_none=True).items()
if k in ("certification_status", "certification_time", "account_type")}
updated = await db.users.set_classification(user_id, **fields)
if updated is None:
raise HTTPException(status_code=404, detail="User not found")
await write_audit(db, action="user.classification", resource="user", resource_id=user_id,
detail=str(fields), user=actor, request=request)
return {"ok": True, "user": updated}
@router.post("/users/{user_id}/status", summary="禁用/启用用户")
async def set_user_status(
user_id: str,
req: SetUserStatusRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.disable")),
):
user = await db.users.get_by_id(user_id)
if user is None:
raise HTTPException(status_code=404, detail="User not found")
updated = await db.users.set_status(user_id, req.status)
# 平台用户启用/禁用 → 引擎镜像同步
await _sync_compute(compute_client.sync_user_enabled(
user.get("username"), req.status == "active",
))
await write_audit(
db, action="user.disable", resource="user", resource_id=user_id,
detail=f"{actor['username']} -> status={req.status}",
user=actor, request=request,
)
return {"ok": True, "user": await db.users.to_profile(updated)}
@router.put("/users/{user_id}", summary="修改用户(资料/角色)")
async def update_user(
user_id: str,
req: UpdateUserRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.manage")),
):
user = await db.users.get_by_id(user_id)
if user is None:
raise HTTPException(status_code=404, detail="User not found")
profile_fields = {k: v for k, v in req.model_dump(exclude_none=True).items() if k in ("nickname",)}
if profile_fields:
user = await db.users.update_profile(user_id, profile_fields) or user
if req.role or req.sub_role is not None:
user = await db.users.set_role(user_id, req.role or user.get("role"), req.sub_role,
req.org_id or user.get("org_id"), req.region_id or user.get("region_id"))
await write_audit(db, action="user.update", resource="user", resource_id=user_id,
detail=str(req.model_dump(exclude_none=True)), user=actor, request=request)
return {"ok": True, "user": await db.users.to_profile(user)}
@router.delete("/users/{user_id}", summary="删除用户")
async def delete_user(
user_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.manage")),
):
user = await db.users.get_by_id(user_id)
ok = await db.users.delete(user_id)
if not ok:
raise HTTPException(status_code=404, detail="User not found")
# 平台用户删除 → 引擎镜像同步删除
if user:
await _sync_compute(compute_client.sync_delete_user(user.get("username")))
await write_audit(db, action="user.delete", resource="user", resource_id=user_id,
detail="deleted", user=actor, request=request)
return {"ok": True, "user_id": user_id}
@router.get("/roles", summary="角色列表")
async def list_roles(
db: Database = Depends(get_db),
_user: dict = Depends(require_permission("menu:admin_role_mgmt")),
):
return await db.roles.list_roles()
@router.get("/permissions", summary="权限列表")
async def list_permissions(
db: Database = Depends(get_db),
_user: dict = Depends(require_permission("menu:admin_role_mgmt")),
):
return await db.roles.list_permissions()
@router.post("/roles/{role_id}/permissions", summary="配置角色权限")
async def set_role_permissions(
role_id: str,
req: RolePermissionRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:role.grant_perm")),
):
if not await db.roles.role_exists(role_id):
raise HTTPException(status_code=404, detail="Role not found")
perms = await db.roles.set_role_permissions(role_id, req.permissions)
await write_audit(db, action="role.grant_perm", resource="role", resource_id=role_id,
detail=f"permissions={len(req.permissions)}", user=actor, request=request)
return {"ok": True, "role_id": role_id, "permissions": perms}
@router.get("/audit-logs", summary="审计日志")
async def list_audit_logs(
limit: int = 100,
offset: int = 0,
db: Database = Depends(get_db),
_user: dict = Depends(require_permission("action:audit.view")),
):
return await db.audit.list(limit=min(limit, 500), offset=offset)
async def _apply_admission_user(db, adm, status):
"""入园/转园审核通过后回填 users 的园区归属。"""
if status not in ("approved", "certified"):
return
uid = (adm or {}).get("user_id")
if not uid:
return
await db.users.set_classification(uid, affiliation="park",
park_id=(adm or {}).get("tenant_id", ""),
park_name=(adm or {}).get("tenant_name", ""))
# ==================== OPC 认证 / 入园 / 转园 审核(平台 operator 兜底) ====================
@router.get("/opc-certifications", summary="OPC 认证申请列表(平台)")
async def list_opc_certifications(
status: str = "",
db: Database = Depends(get_db),
_role: dict = Depends(require_roles("operator")),
_perm: dict = Depends(require_permission("menu:admin_user_mgmt")),
):
return {"items": await db.opc_certifications.list(status or None)}
@router.post("/opc-certifications/{cert_id}/review", summary="审核 OPC 认证申请")
async def review_opc_certification(
cert_id: str,
req: OpcCertReviewRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.manage")),
):
cert = await db.opc_certifications.get(cert_id)
if cert is None:
raise HTTPException(status_code=404, detail="认证申请不存在")
if req.status not in ("certified", "rejected"):
raise HTTPException(status_code=400, detail="状态仅支持 certified/rejected")
updated = await db.opc_certifications.set_status(cert_id, req.status,
reviewer=actor.get("id", ""), comment=req.comment)
if cert.get("user_id"):
u = await db.users.get_by_id(cert["user_id"])
await db.users.set_classification(cert["user_id"],
certification_status=req.status, affiliation=(u or {}).get("affiliation") or "independent")
await write_audit(db, action="opc.cert_review", resource="opc_certification",
resource_id=cert_id, detail=f"status={req.status}", user=actor, request=request)
return {"ok": True, "certification": updated}
@router.get("/park-admissions", summary="园区入驻申请列表(平台)")
async def list_park_admissions(
status: str = "",
db: Database = Depends(get_db),
_role: dict = Depends(require_roles("operator")),
_perm: dict = Depends(require_permission("menu:admin_user_mgmt")),
):
return {"items": await db.park_admissions.list(status or None)}
@router.post("/park-admissions/{aid}/review", summary="园区入驻申请审核(平台兜底)")
async def review_park_admission(
aid: str,
req: ParkAdmissionReviewRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.manage")),
):
adm = await db.park_admissions.get(aid)
if adm is None:
raise HTTPException(status_code=404, detail="入驻申请不存在")
if req.status not in ("approved", "rejected", "reviewing"):
raise HTTPException(status_code=400, detail="状态仅支持 approved/rejected/reviewing")
updated = await db.park_admissions.set_status(aid, req.status,
reviewer=actor.get("id", ""), comment=req.comment)
await _apply_admission_user(db, adm, req.status)
await write_audit(db, action="opc.admission_review", resource="park_admission",
resource_id=aid, detail=f"status={req.status}", user=actor, request=request)
return {"ok": True, "admission": updated}
@router.get("/park-transfers", summary="OPC 转园申请列表(平台)")
async def list_park_transfers(
status: str = "",
db: Database = Depends(get_db),
_role: dict = Depends(require_roles("operator")),
_perm: dict = Depends(require_permission("menu:admin_user_mgmt")),
):
return {"items": await db.park_transfers.list(status or None)}
@router.post("/park-transfers/{tid}/review", summary="OPC 转园申请审核(平台兜底)")
async def review_park_transfer(
tid: str,
req: ParkTransferReviewRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.manage")),
):
tr = await db.park_transfers.get(tid)
if tr is None:
raise HTTPException(status_code=404, detail="转园申请不存在")
if req.status not in ("approved", "rejected", "reviewing"):
raise HTTPException(status_code=400, detail="状态仅支持 approved/rejected/reviewing")
updated = await db.park_transfers.set_status(tid, req.status,
reviewer=actor.get("id", ""), comment=req.comment)
if req.status == "approved":
await db.users.set_classification(tr["user_id"], affiliation="park",
park_id=tr.get("to_park_id", ""), park_name=tr.get("to_park_name", ""))
await write_audit(db, action="opc.park_transfer_review", resource="park_transfer",
resource_id=tid, detail=f"status={req.status}", user=actor, request=request)
return {"ok": True, "transfer": updated}
# ==================== 平台园区管理(operator;复用 park 数据层,不依赖 app.state.db ====================
class AdminBindDeviceBody(BaseModel):
tenant_id: str
name: str = ""
role: str = "main"
location: str = ""
class AdminCreateTenantBody(BaseModel):
name: str
intro: list[str] = ["", ""]
username: str = ""
password: str = ""
class AdminTenantStatusBody(BaseModel):
status: str # active | disabled
def _park_tenant_public(t: dict) -> dict:
"""脱敏园区:剔除 auth(username/salt/password_hash),仅留管理用字段。"""
return {k: t.get(k) for k in ("id", "name", "intro", "admin_username", "status")}
@router.get("/park/screens", summary="全部屏幕(含未绑定+归属园区名,平台)")
async def admin_park_screens(
db: Database = Depends(get_db),
_role: dict = Depends(require_roles("operator")),
_perm: dict = Depends(require_permission("menu:admin_user_mgmt")),
):
from app.park import tenants as tnt
return {"items": await tnt.list_all_screens()}
@router.get("/park/tenants", summary="园区列表(脱敏,平台)")
async def admin_park_tenants(
db: Database = Depends(get_db),
_role: dict = Depends(require_roles("operator")),
_perm: dict = Depends(require_permission("menu:admin_user_mgmt")),
):
from app.park import tenants as tnt
rows = await tnt.list_tenants()
return {"items": [_park_tenant_public(t) for t in rows]}
@router.post("/park/tenants", summary="平台创建园区")
async def admin_park_create_tenant(
body: AdminCreateTenantBody,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.manage")),
):
from app.park import tenants as tnt
t = await tnt.create_tenant(body.name, body.intro, body.username, body.password)
await write_audit(db, action="park.create", resource="park_tenant",
resource_id=(t or {}).get("id", ""), detail=body.name, user=actor, request=request)
return {"ok": True, "tenant": _park_tenant_public(t or {})}
@router.delete("/park/tenants/{tid}", summary="平台删除园区(级联删屏)")
async def admin_park_delete_tenant(
tid: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.manage")),
):
from app.park import tenants as tnt
ok = await tnt.delete_tenant(tid)
if not ok:
raise HTTPException(status_code=404, detail="园区不存在")
await write_audit(db, action="park.delete", resource="park_tenant",
resource_id=tid, detail=tid, user=actor, request=request)
return {"ok": True}
@router.post("/park/tenants/{tid}/status", summary="平台停用/启用园区")
async def admin_park_tenant_status(
tid: str,
body: AdminTenantStatusBody,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.manage")),
):
from app.park import tenants as tnt
if body.status not in ("active", "disabled"):
raise HTTPException(status_code=400, detail="无效状态")
ok = await tnt.set_status(tid, body.status)
if not ok:
raise HTTPException(status_code=404, detail="园区不存在")
await write_audit(db, action="park.status", resource="park_tenant",
resource_id=tid, detail=body.status, user=actor, request=request)
return {"ok": True, "status": body.status}
@router.post("/park/devices/{device_id}/bind", summary="平台手动绑定屏幕到园区")
async def admin_park_bind_device(
device_id: str,
body: AdminBindDeviceBody,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.manage")),
):
from app.park import tenants as tnt
tenant = await tnt.get_tenant(body.tenant_id)
if tenant is None:
raise HTTPException(status_code=404, detail="目标园区不存在")
dev = await tnt.get_device_by_id(device_id)
if dev is None:
raise HTTPException(status_code=404, detail="屏幕不存在")
if dev.get("status") == "bound" and dev.get("tenant_id") and dev.get("tenant_id") != body.tenant_id:
raise HTTPException(status_code=400, detail="该屏幕已绑定到其它园区,请先解绑")
updated = await tnt.bind_device(body.tenant_id, device_id, body.name, body.role, body.location)
await write_audit(db, action="park.bind_screen", resource="park_screen",
resource_id=device_id, detail=f"tid={body.tenant_id}", user=actor, request=request)
return {"ok": True, "screen": updated}
@router.post("/park/devices/{device_id}/unbind", summary="平台解绑屏幕(从园区解绑,回未绑定)")
async def admin_park_unbind_device(
device_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.manage")),
):
from app.park import tenants as tnt
dev = await tnt.get_device_by_id(device_id)
if dev is None:
raise HTTPException(status_code=404, detail="屏幕不存在")
if dev.get("status") != "bound":
raise HTTPException(status_code=400, detail="该屏幕当前未绑定,无需解绑")
ok = await tnt.unbind_device(device_id)
if not ok:
raise HTTPException(status_code=500, detail="解绑失败")
await write_audit(db, action="park.unbind_screen", resource="park_screen",
resource_id=device_id, detail=f"from tid={dev.get('tenant_id')}", user=actor, request=request)
return {"ok": True}
class AdminCompanyComputeBody(BaseModel):
discount: int | None = Field(default=None, ge=0, le=100)
quota_add: int = Field(default=0, ge=0)
@router.get("/park/tenants/{tid}/companies", summary="园区企业列表(含成员数/算力,平台)")
async def admin_park_tenant_companies(
tid: str,
db: Database = Depends(get_db),
_role: dict = Depends(require_roles("operator")),
_perm: dict = Depends(require_permission("menu:admin_user_mgmt")),
):
from app.park import tenants as tnt
return {"items": await tnt.list_companies(tid)}
@router.put("/park/tenants/{tid}/companies/{cid}/compute", summary="设置企业算力(只增不减,平台兜底)")
async def admin_park_tenant_company_compute(
tid: str,
cid: str,
body: AdminCompanyComputeBody,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:user.manage")),
):
from app.park import tenants as tnt
c = await tnt.get_company(cid)
if c is None or c.get("tenant_id") != tid:
raise HTTPException(status_code=404, detail="企业不存在")
try:
updated = await tnt.set_company_compute(cid, discount=body.discount, quota_add=body.quota_add)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc))
await tnt.sync_company_engine(cid)
if body.quota_add > 0:
from ..services import compute_client
for m in await tnt.company_members(cid):
try:
await compute_client.grant_user_quota_by_username(m["username"], body.quota_add)
except Exception: # noqa: BLE001
continue
await write_audit(db, action="park.company_compute", resource="park_company",
resource_id=cid, detail=f"discount={body.discount} quota_add={body.quota_add}", user=actor, request=request)
return {"ok": True, "company": updated, "quota_add": body.quota_add}