feat(opc): OPC认证/入园/转园 端点 + 园区端方案B审核 + 迁移0013

- rbac_opc.py: /opc/certification/apply|mine、/opc/park-transfer/apply|mine
- rbac_admin.py: platform operator 兜底 list/review (opc-certifications/park-admissions/park-transfers),通过回填 users.set_classification
- park/routers.py: carrier 角色凭 operator_user_id 定位本园区,/park/api park-admissions|transfers view+review
- park/tenants.py: bind_admin 写 operator_user_id + find_by_operator_user_id
- 迁移 0013: opc_certifications + park_transfers + park_tenants.operator_user_id
This commit is contained in:
Pine
2026-08-26 15:43:34 +08:00
parent 30f0baf074
commit bc6265eb67
5 changed files with 377 additions and 3 deletions
+99 -1
View File
@@ -21,6 +21,8 @@ from sqlalchemy.ext.asyncio import AsyncSession
from ..infrastructure.db import get_session
from ..infrastructure.repositories import TaskRepository
from ..rbac import require_roles, write_audit
from ..api.schemas.admin import ParkAdmissionReviewRequest, ParkTransferReviewRequest
from . import park_config, tenants
from .auth import create_device_token, create_token, parse_device_token, parse_token, require_tenant
from .config import settings
@@ -140,11 +142,12 @@ async def tenant_bind(tid: str, body: BindBody):
db = Database()
try:
user = await db.users.get_by_username(body.username.strip())
user_id = (user or {}).get("id", "")
finally:
await db.close()
if user is None:
return JSONResponse({"ok": False, "error": "账号不存在"}, status_code=404)
ok = await tenants.bind_admin(tid, body.username.strip())
ok = await tenants.bind_admin(tid, body.username.strip(), operator_user_id=user_id)
return {"ok": ok, "tenant_id": tid, "admin_username": body.username.strip()}
@@ -864,3 +867,98 @@ def asyncio_timeout(awaitable, seconds):
"""极简超时包装:避免依赖 asyncio.timeout3.11+ 的上下文管理器)。"""
import asyncio
return asyncio.wait_for(awaitable, timeout=seconds)
# ==================== 园区端(carrier)入园 / 转园 审核 —— 仅本园区 ====================
async def _carrier_db():
"""园区端审核端点统一走平台 Database(与 park_login 一致)。"""
from app.infrastructure.repositories import Database
db = Database()
try:
yield db
finally:
await db.close()
async def _my_park(user: dict) -> dict:
"""当前 carrier 账号关联的园区(operator_user_id),未关联则拒。"""
t = await tenants.find_by_operator_user_id(user["id"])
if not t:
raise HTTPException(status_code=403, detail="该账号未关联任何园区")
return t
@router.get("/api/park-admissions", summary="园区入驻申请(本园区, carrier)")
async def carrier_park_admissions(
status: str = "",
db=Depends(_carrier_db),
user: dict = Depends(require_roles("carrier")),
):
t = await _my_park(user)
return {"items": await db.park_admissions.list(status or None, tenant_id=t["id"])}
@router.post("/api/park-admissions/{aid}/review", summary="园区入驻申请审核(本园区)")
async def carrier_park_admission_review(
aid: str,
body: ParkAdmissionReviewRequest,
request: Request,
db=Depends(_carrier_db),
user: dict = Depends(require_roles("carrier")),
):
t = await _my_park(user)
adm = await db.park_admissions.get(aid)
if adm is None:
raise HTTPException(status_code=404, detail="入驻申请不存在")
if adm.get("tenant_id") != t["id"]:
raise HTTPException(status_code=403, detail="仅可审核本园区申请")
st = body.status
if st not in ("approved", "rejected", "reviewing"):
raise HTTPException(status_code=400, detail="状态仅支持 approved/rejected/reviewing")
updated = await db.park_admissions.set_status(aid, st,
reviewer=user.get("id", ""), comment=body.comment)
if st == "approved" and adm.get("user_id"):
await db.users.set_classification(adm["user_id"], affiliation="park",
park_id=adm.get("tenant_id", ""), park_name=adm.get("tenant_name", ""))
await write_audit(db, action="park.admission_review", resource="park_admission",
resource_id=aid, detail=f"status={st}", user=user, request=request)
return {"ok": True, "admission": updated}
@router.get("/api/park-transfers", summary="OPC 转园申请(本园区相关, carrier)")
async def carrier_park_transfers(
status: str = "",
db=Depends(_carrier_db),
user: dict = Depends(require_roles("carrier")),
):
t = await _my_park(user)
return {"items": await db.park_transfers.list(status or None, park_id=t["id"])}
@router.post("/api/park-transfers/{tid}/review", summary="OPC 转园申请审核(本园区)")
async def carrier_park_transfer_review(
tid: str,
body: ParkTransferReviewRequest,
request: Request,
db=Depends(_carrier_db),
user: dict = Depends(require_roles("carrier")),
):
t = await _my_park(user)
tr = await db.park_transfers.get(tid)
if tr is None:
raise HTTPException(status_code=404, detail="转园申请不存在")
if tr.get("from_park_id") != t["id"] and tr.get("to_park_id") != t["id"]:
raise HTTPException(status_code=403, detail="仅可审核与本园区相关的转园")
st = body.status
if st not in ("approved", "rejected", "reviewing"):
raise HTTPException(status_code=400, detail="状态仅支持 approved/rejected/reviewing")
updated = await db.park_transfers.set_status(tid, st,
reviewer=user.get("id", ""), comment=body.comment)
if st == "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="park.transfer_review", resource="park_transfer",
resource_id=tid, detail=f"status={st}", user=user, request=request)
return {"ok": True, "transfer": updated}
+10 -2
View File
@@ -205,9 +205,10 @@ async def _append_companies(s, tenant_id: str, companies: list) -> None:
founder=c.get("founder", "") if isinstance(c, dict) else "", status="active", employees=None, created_at=_now()))
async def bind_admin(tenant_id: str, username: str) -> bool:
async def bind_admin(tenant_id: str, username: str, operator_user_id: str = "") -> bool:
async with _get_session() as s:
r = await s.execute(update(ParkTenant).where(ParkTenant.id == tenant_id).values(admin_username=username.strip()))
r = await s.execute(update(ParkTenant).where(ParkTenant.id == tenant_id)
.values(admin_username=username.strip(), operator_user_id=operator_user_id or ""))
await s.commit()
return r.rowcount > 0
@@ -222,6 +223,13 @@ async def find_by_admin(username: str) -> dict | None:
return _to_tenant(row) if row else None
async def find_by_operator_user_id(user_id: str) -> dict | None:
"""carrier 载体账号据此定位本园区(园区端只审本片区)。"""
async with _get_session() as s:
row = (await s.execute(select(ParkTenant).where(ParkTenant.operator_user_id == user_id).limit(1))).scalars().first()
return _to_tenant(row) if row else None
async def update_tenant(tenant_id: str, patch: dict) -> dict | None:
async with _get_session() as s:
row = await s.get(ParkTenant, tenant_id)