diff --git a/alembic/versions/0013_opc_cert_park_transfer.py b/alembic/versions/0013_opc_cert_park_transfer.py new file mode 100644 index 0000000..66aa603 --- /dev/null +++ b/alembic/versions/0013_opc_cert_park_transfer.py @@ -0,0 +1,66 @@ +"""opc_certifications + park_transfers + park_tenants 归属列 + +Revision ID: 0013_opc_cert_park_transfer +Revises: 0012_park_admissions +Create Date: 2026-08-26 +""" +from __future__ import annotations + +import sqlalchemy as sa +from alembic import op + +revision = "0013_opc_cert_park_transfer" +down_revision = "0012_park_admissions" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + # OPC 认证申请表(用户提交资料 → 运营方审核) + op.create_table( + "opc_certifications", + sa.Column("id", sa.String(), primary_key=True), + sa.Column("user_id", sa.String(), default="", index=True), + sa.Column("username", sa.String(), default=""), + sa.Column("real_name", sa.String(), default=""), + sa.Column("gender", sa.String(), default=""), + sa.Column("address", sa.String(), default=""), + sa.Column("industry", sa.String(), default=""), + sa.Column("ability", sa.Text(), default=""), + sa.Column("phone", sa.String(), default=""), + sa.Column("docs_json", sa.Text(), default="{}"), + sa.Column("status", sa.String(), default="pending"), + sa.Column("review_comment", sa.String(), default=""), + sa.Column("reviewed_by", sa.String(), default=""), + sa.Column("reviewed_at", sa.String(), default=""), + sa.Column("created_at", sa.String(), default=""), + sa.Column("updated_at", sa.String(), default=""), + ) + + # OPC 转园申请表(园区 OPC 发起 → 园区/平台审核) + op.create_table( + "park_transfers", + sa.Column("id", sa.String(), primary_key=True), + sa.Column("user_id", sa.String(), default="", index=True), + sa.Column("username", sa.String(), default=""), + sa.Column("from_park_id", sa.String(), default=""), + sa.Column("from_park_name", sa.String(), default=""), + sa.Column("to_park_id", sa.String(), default=""), + sa.Column("to_park_name", sa.String(), default=""), + sa.Column("reason", sa.Text(), default=""), + sa.Column("status", sa.String(), default="pending"), + sa.Column("review_comment", sa.String(), default=""), + sa.Column("reviewed_by", sa.String(), default=""), + sa.Column("reviewed_at", sa.String(), default=""), + sa.Column("created_at", sa.String(), default=""), + sa.Column("updated_at", sa.String(), default=""), + ) + + # 园区→载体归属:园区端据此只审本园区 + op.add_column("park_tenants", sa.Column("operator_user_id", sa.String(), default="")) + + +def downgrade() -> None: + op.drop_column("park_tenants", "operator_user_id") + op.drop_table("park_transfers") + op.drop_table("opc_certifications") diff --git a/app/api/routers/rbac_admin.py b/app/api/routers/rbac_admin.py index 772ebfb..453a93e 100644 --- a/app/api/routers/rbac_admin.py +++ b/app/api/routers/rbac_admin.py @@ -8,6 +8,7 @@ from pydantic import BaseModel 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 @@ -204,3 +205,115 @@ async def list_audit_logs( _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} diff --git a/app/api/routers/rbac_opc.py b/app/api/routers/rbac_opc.py index 7db390f..007d68b 100644 --- a/app/api/routers/rbac_opc.py +++ b/app/api/routers/rbac_opc.py @@ -386,3 +386,92 @@ async def opc_compute_tokens_delete(token_id: int, user: dict = Depends(require_ except compute_client.ComputeError as exc: raise HTTPException(status_code=502, detail=f"算力引擎对接失败: {exc}") from exc return {"ok": True, "token_id": token_id} + + +# ── OPC 认证 ───────────────────────────────────────────────────────────── + +@router.post("/certification/apply", summary="提交 OPC 认证申请") +async def opc_cert_apply( + req: OpcCertApplyRequest, + request: Request, + db: Database = Depends(get_db), + user: dict = Depends(require_roles("opc_member")), +): + """当前 OPC 提交认证资料 → 状态置 pending,待运营方审核。""" + me = await db.users.get_by_id(user["id"]) + if me is None: + raise HTTPException(status_code=404, detail="用户不存在") + if me.get("certification_status") == "certified": + raise HTTPException(status_code=400, detail="当前账号已是认证OPC,无需重复申请") + if me.get("certification_status") == "pending": + raise HTTPException(status_code=400, detail="认证申请审核中,请勿重复提交") + applied = await db.opc_certifications.by_user(user["id"]) + if applied and applied.get("status") in ("pending", "reviewing"): + raise HTTPException(status_code=400, detail="已有待审核的认证申请,请勿重复提交") + cert = await db.opc_certifications.create({ + "real_name": req.real_name, "gender": req.gender, "address": req.address, + "industry": req.industry, "ability": req.ability, "phone": req.phone, + "docs_json": req.docs_json or "{}", + }, user["id"], me.get("username", "")) + await write_audit( + action="opc.cert_apply", resource="opc_certification", resource_id=cert["id"], + detail=f"real_name={req.real_name} industry={req.industry}", + user=user, request=request, + ) + return {"ok": True, "certification": cert} + + +@router.get("/certification/mine", summary="我的 OPC 认证状态") +async def opc_cert_mine( + db: Database = Depends(get_db), + user: dict = Depends(require_roles("opc_member")), +): + me = await db.users.get_by_id(user["id"]) + cert = await db.opc_certifications.by_user(user["id"]) + return { + "certification_status": (me or {}).get("certification_status", "uncertified"), + "certification": cert, + } + + +# ── OPC 转园 ───────────────────────────────────────────────────────────── + +@router.post("/park-transfer/apply", summary="发起 OPC 转园申请") +async def opc_park_transfer_apply( + req: ParkTransferApplyRequest, + request: Request, + db: Database = Depends(get_db), + user: dict = Depends(require_roles("opc_member")), +): + """园区 OPC 申请转入目标园区(原园区自动取当前 park 归属,经园区/平台审核)。""" + me = await db.users.get_by_id(user["id"]) + if me is None: + raise HTTPException(status_code=404, detail="用户不存在") + if me.get("affiliation") != "park" or not me.get("park_id"): + raise HTTPException(status_code=400, detail="仅园区 OPC 可发起转园(需先归属园区)") + if not req.to_park_id: + raise HTTPException(status_code=400, detail="请选择目标园区") + if req.to_park_id == me.get("park_id"): + raise HTTPException(status_code=400, detail="目标园区与当前园区相同") + existing = await db.park_transfers.by_user(user["id"]) + if existing and existing.get("status") in ("pending", "reviewing"): + raise HTTPException(status_code=400, detail="已有待审核的转园申请,请勿重复提交") + transfer = await db.park_transfers.create({ + "to_park_id": req.to_park_id, "to_park_name": req.to_park_name, + "reason": req.reason, + "from_park_id": me.get("park_id", ""), "from_park_name": me.get("park_name", ""), + }, user["id"], me.get("username", "")) + await write_audit( + action="opc.park_transfer_apply", resource="park_transfer", resource_id=transfer["id"], + detail=f"from={me.get('park_name','')} to={req.to_park_name}", + user=user, request=request, + ) + return {"ok": True, "transfer": transfer} + + +@router.get("/park-transfer/mine", summary="我的转园申请") +async def opc_park_transfer_mine( + db: Database = Depends(get_db), + user: dict = Depends(require_roles("opc_member")), +): + return {"transfer": await db.park_transfers.by_user(user["id"])} diff --git a/app/park/routers.py b/app/park/routers.py index 52c833b..02f0ae0 100644 --- a/app/park/routers.py +++ b/app/park/routers.py @@ -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.timeout(3.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} + diff --git a/app/park/tenants.py b/app/park/tenants.py index 6087e7b..a48ba4a 100644 --- a/app/park/tenants.py +++ b/app/park/tenants.py @@ -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)