# -*- coding: utf-8 -*- """运营端业务管理端点:任务 / 服务商 / 内容 / 系统配置 / 数据统计。 全部要求 operator 角色(内部子角色再按权限细分),并写审计日志。 """ from __future__ import annotations import json from fastapi import APIRouter, Depends, HTTPException, Request, Response from pydantic import BaseModel from ..dependencies import get_db from ..schemas.operator import TaskCreateRequest, TaskStatusRequest, ProviderCreateRequest, ProviderUpdateRequest, ContentCreateRequest, ContentStatusRequest, ConfigUpdateRequest, ComputePingResponse, ComputeProvisionRequest, ComputeProvisionResponse, CourseCreateRequest, CourseStatusRequest, ActivityStatusRequest, BookingUpdateRequest, TestCreateRequest, TestStatusRequest from ...rbac import require_permission, require_roles, write_audit from ...infrastructure.repositories import Database from ...services import compute_client from ...services import training_admin_bridge as training router = APIRouter(prefix="/admin", tags=["admin-op"]) @router.get("/tasks", summary="任务列表") async def list_tasks( status: str | None = None, db: Database = Depends(get_db), _u: dict = Depends(require_roles("operator")), ): return await db.tasks.list(status=status) @router.post("/tasks", summary="创建任务") async def create_task( req: TaskCreateRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_permission("action:task.manage")), ): task = await db.tasks.create(req.model_dump(exclude_none=True)) await write_audit(db, action="task.create", resource="task", resource_id=task["id"], detail=task["title"], user=actor, request=request) return task @router.post("/tasks/{task_id}/status", summary="更新任务状态") async def set_task_status( task_id: str, req: TaskStatusRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_permission("action:task.manage")), ): if await db.tasks.get(task_id) is None: raise HTTPException(status_code=404, detail="Task not found") task = await db.tasks.set_status(task_id, req.status) await write_audit(db, action="task.status", resource="task", resource_id=task_id, detail=req.status, user=actor, request=request) return task # ── 服务商管理 ──────────────────────────────────────────────────────────── @router.get("/providers", summary="服务商列表") async def list_providers( status: str | None = None, db: Database = Depends(get_db), _u: dict = Depends(require_roles("operator")), ): return await db.providers.list(status=status) @router.post("/providers", summary="创建服务商") async def create_provider( req: ProviderCreateRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_permission("action:provider.manage")), ): prov = await db.providers.create(req.model_dump(exclude_none=True)) await write_audit(db, action="provider.create", resource="provider", resource_id=prov["id"], detail=prov["name"], user=actor, request=request) return prov @router.post("/providers/{provider_id}/update", summary="更新服务商(评级/状态)") async def update_provider( provider_id: str, req: ProviderUpdateRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_permission("action:provider.manage")), ): updated = await db.providers.update(provider_id, req.model_dump(exclude_none=True)) if updated is None: raise HTTPException(status_code=404, detail="Provider not found") await write_audit(db, action="provider.update", resource="provider", resource_id=provider_id, detail=str(req.model_dump(exclude_none=True)), user=actor, request=request) return updated # ── 内容管理 ────────────────────────────────────────────────────────────── @router.get("/content", summary="内容列表") async def list_content( ctype: str | None = None, status: str | None = None, db: Database = Depends(get_db), _u: dict = Depends(require_roles("operator")), ): return await db.content.list(ctype=ctype, status=status) @router.post("/content", summary="创建内容") async def create_content( req: ContentCreateRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_permission("action:content.manage")), ): item = await db.content.create(req.model_dump(exclude_none=True)) await write_audit(db, action="content.create", resource="content", resource_id=item["id"], detail=item["title"], user=actor, request=request) return item @router.post("/content/{content_id}/status", summary="发布/下架内容") async def set_content_status( content_id: str, req: ContentStatusRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_permission("action:content.manage")), ): item = await db.content.set_status(content_id, req.status) if item is None: raise HTTPException(status_code=404, detail="Content not found") await write_audit(db, action="content.status", resource="content", resource_id=content_id, detail=req.status, user=actor, request=request) return item # ── 系统配置 ────────────────────────────────────────────────────────────── @router.get("/config", summary="系统配置列表") async def list_config( db: Database = Depends(get_db), _u: dict = Depends(require_roles("operator")), ): return await db.config.all() @router.put("/config/{key}", summary="更新系统配置") async def set_config( key: str, req: ConfigUpdateRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_permission("action:config.manage")), ): cfg = await db.config.set(key, req.value, req.description) await write_audit(db, action="config.update", resource="config", resource_id=key, detail=req.value, user=actor, request=request) return cfg # ── 数据统计 ────────────────────────────────────────────────────────────── @router.get("/stats/overview", summary="运营数据总览") async def stats_overview( db: Database = Depends(get_db), _u: dict = Depends(require_roles("operator")), ): return await db.stats.overview() # ── 算力对接(compute-engine) ───────────────────────────────────────────── @router.post("/compute/ping", response_model=ComputePingResponse, summary="算力引擎连通性检查") async def compute_ping( db: Database = Depends(get_db), _u: dict = Depends(require_roles("operator")), ): data = await compute_client.ping() status = str(data.get("status", "")).lower() return ComputePingResponse( ok=status in ("ok", "success") or "success" in status or not data.get("detail"), version=data.get("version"), detail=data.get("detail"), ) @router.post("/compute/provision", response_model=ComputeProvisionResponse, summary="开通算力(建号+发PAT)") async def compute_provision( req: ComputeProvisionRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_roles("operator")), ): """为平台用户开通算力:引擎建号 + 签发 PAT,返回给前端注入 provider。""" user = await db.users.get_by_id(req.user_id) if user is None: raise HTTPException(status_code=404, detail="用户不存在") username = user.get("username") or user["id"] try: await compute_client.create_user(username) pat = await compute_client.issue_pat(username) created = True except compute_client.ComputeError as exc: raise HTTPException(status_code=502, detail=f"算力引擎对接失败: {exc}") from exc await write_audit(db, action="compute.provision", resource="compute", resource_id=req.user_id, detail=f"engine user {username} provisioned", user=actor, request=request) return ComputeProvisionResponse( engine_user_id=username, engine_username=username, pat=pat, created=created, ) @router.post("/compute/sync-users", summary="同步引擎用户=平台总用户(对账)") async def sync_compute_users( request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_roles("operator")), ): """把平台全部用户对账到引擎:建号(幂等) + 缺令牌则签一枚 PAT,保证 engine 用户 = 平台用户。""" users = await db.users.list() created = pats = 0 for u in users: uname = u.get("username") if not uname: continue try: await compute_client.ensure_user(uname) created += 1 except Exception: # noqa: BLE001 pass try: toks = await compute_client.list_user_tokens(uname) if not toks: await compute_client.issue_pat(uname) pats += 1 except Exception: # noqa: BLE001 pass # 回写算力引擎镜像(provisioned/quota/used)到平台用户,供直观展示 try: bal = await compute_client.user_balance(uname) await db.users.set_compute_mirror( u["id"], provisioned=True, username=uname, quota=int(bal.get("quota", 0) or 0), used_quota=int(bal.get("used_quota", 0) or 0), ) except Exception: # noqa: BLE001 pass # 对账幂等自明、且涉及对 db 只读 + 大量外部调用,不写审计(write_audit 的 commit 与该只读会话 # 事务交互会触发 PendingRollback/UNIQUE 冲突并污染会话)。仅返回统计。 return {"total": len(users), "synced": created, "pat_issued": pats} @router.api_route( "/compute/proxy/{path:path}", methods=["GET", "POST", "PUT", "DELETE", "PATCH"], response_class=Response, summary="算力中心管理 API 透传(对齐 compute-engine /api/*)", ) async def compute_proxy(path: str, request: Request, _u: dict = Depends(require_roles("operator"))): """算力中心管理转发:把 admin 端请求 1:1 透传到 compute-engine(new-api) 管理 API ``/api/{path}``。 打通全部管理面(models / channels / groups / tokens / users / logs / redemption / ratio / system-settings 等),保留引擎原始状态码与响应体,使 admin 端可作为 new-api 唯一管理 UI。 """ body = await request.body() json_body = None if body and request.headers.get("content-type", "").startswith("application/json"): try: json_body = json.loads(body) except Exception: # noqa: BLE001 json_body = None params = dict(request.query_params) result = await compute_client.proxy( request.method, "/" + path, json_body=json_body, params=params, ) return Response(content=result["body"], status_code=int(result["status"]), media_type="application/json") # ── 培训业务 · 课程(桥接 app.training courses 表)───────────────────────── @router.get("/courses", summary="课程列表") async def list_courses( status: str | None = None, _u: dict = Depends(require_roles("operator")), ): return training.list_courses(status=status) @router.post("/courses", summary="创建课程") async def create_course( req: CourseCreateRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_permission("action:course.manage")), ): course = training.create_course(req.model_dump(exclude_none=True)) await write_audit(db, action="course.create", resource="course", resource_id=course["id"], detail=course["title"], user=actor, request=request) return course @router.post("/courses/{course_id}/status", summary="更新课程状态") async def set_course_status( course_id: str, req: CourseStatusRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_permission("action:course.manage")), ): course = training.set_course_status(course_id, req.status) if course is None: raise HTTPException(status_code=404, detail="Course not found") await write_audit(db, action="course.status", resource="course", resource_id=course_id, detail=req.status, user=actor, request=request) return course # ── 培训业务 · 活动(桥接 app.training events 表)────────────────────────── @router.get("/activities", summary="活动列表") async def list_activities( status: str | None = None, _u: dict = Depends(require_roles("operator")), ): return training.list_activities(status=status) @router.post("/activities/{event_id}/status", summary="更新活动状态") async def set_activity_status( event_id: str, req: ActivityStatusRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_permission("action:activity.manage")), ): activity = training.set_activity_status(event_id, req.status) if activity is None: raise HTTPException(status_code=404, detail="Activity not found") await write_audit(db, action="activity.status", resource="activity", resource_id=event_id, detail=req.status, user=actor, request=request) return activity # ── 培训业务 · 报名(桥接 app.training bookings 表)──────────────────────── @router.get("/bookings", summary="报名列表") async def list_bookings( status: str | None = None, audit_status: str | None = None, _u: dict = Depends(require_roles("operator")), ): return training.list_bookings(status=status, audit_status=audit_status) @router.patch("/bookings/{booking_id}", summary="更新报名(审核/备注)") async def update_booking( booking_id: str, req: BookingUpdateRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_permission("action:booking.manage")), ): booking = training.update_booking(booking_id, req.model_dump(exclude_none=True)) if booking is None: raise HTTPException(status_code=404, detail="Booking not found") await write_audit(db, action="booking.update", resource="booking", resource_id=booking_id, detail=str(req.model_dump(exclude_none=True)), user=actor, request=request) return booking # ── 培训业务 · 测评(桥接 app.training tests 表)────────────────────────── @router.get("/tests", summary="测评记录列表") async def list_tests( _u: dict = Depends(require_roles("operator")), ): return training.list_tests() @router.post("/tests", summary="创建测评记录") async def create_test( req: TestCreateRequest, request: Request, db: Database = Depends(get_db), actor: dict = Depends(require_permission("action:test.manage")), ): item = training.create_test(req.model_dump(exclude_none=True)) await write_audit(db, action="test.create", resource="test", resource_id=item["id"], detail=item["title"], user=actor, request=request) return item @router.post("/tests/{test_id}/status", summary="更新测评状态") async def set_test_status( test_id: str, req: TestStatusRequest, _u: dict = Depends(require_roles("operator")), ): # 测评记录为不可变快照,状态变更仅回读当前记录(后续如需启用/禁用再扩展字段) item = training.pack_test_if_exists(test_id) if item is None: raise HTTPException(status_code=404, detail="Test not found") return item