b745328d49
- 运营端 /admin/tasks/{id} 详情、/{id}/bids、/{id}/claims、/{id}/bids/{bid_id}/win 开标定标(余标置lost+进 in_progress)
- 园区端 POST /park/api/park/tasks 发布任务(publisher_role=park/park_id=本园/发布即上架/默认公有)
817 lines
32 KiB
Python
817 lines
32 KiB
Python
# -*- coding: utf-8 -*-
|
||
"""运营端业务管理端点:任务 / 服务商 / 内容 / 系统配置 / 数据统计。
|
||
|
||
全部要求 operator 角色(内部子角色再按权限细分),并写审计日志。
|
||
"""
|
||
from __future__ import annotations
|
||
|
||
import json
|
||
|
||
from fastapi import APIRouter, Depends, File, HTTPException, Request, Response, UploadFile
|
||
from pydantic import BaseModel
|
||
|
||
from ..dependencies import get_db
|
||
from ..schemas.operator import TaskCreateRequest, TaskUpdateRequest, TaskStatusRequest, TaskAssignRequest, TaskRecommendRequest, TaskSelectRecommendRequest, ProviderCreateRequest, ProviderUpdateRequest, ContentCreateRequest, ContentStatusRequest, ContentUpdateRequest, ContentReviewRequest, ConfigUpdateRequest, ComputePingResponse, ComputeBalanceRequest, ComputeProvisionRequest, ComputeProvisionResponse, CourseCreateRequest, CourseStatusRequest, CourseUpdateRequest, ActivityCreateRequest, ActivityStatusRequest, ActivityUpdateRequest, 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"])
|
||
|
||
|
||
class EscrowDepositRequest(BaseModel):
|
||
amount: int = 0
|
||
|
||
|
||
@router.get("/escrows", summary="运营结算台账(资金托管)")
|
||
async def operator_escrows(
|
||
status: str | None = None,
|
||
db: Database = Depends(get_db),
|
||
_u: dict = Depends(require_roles("operator")),
|
||
):
|
||
"""运营平台资金托管台账:全部任务托管、佣金、状态(deposited/frozen/released/refunded)。"""
|
||
return {"items": await db.escrows.list(status=status)}
|
||
|
||
|
||
@router.post("/tasks/{task_id}/escrow", summary="任务托管建档/充值(运营平台托管)")
|
||
async def operator_escrow_deposit(
|
||
task_id: str,
|
||
req: EscrowDepositRequest,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
_u: dict = Depends(require_roles("operator")),
|
||
):
|
||
"""发标方向运营平台托管账户建档/充值;已托管返回现有。金额默认取任务预算。"""
|
||
from ...services.settlement_service import COMMISSION_RATE
|
||
|
||
task = await db.tasks.get(task_id)
|
||
if task is None:
|
||
raise HTTPException(status_code=404, detail="任务不存在")
|
||
for e in await db.escrows.list():
|
||
if e["task_id"] == task_id:
|
||
return e
|
||
amount = req.amount or max(task.get("budget_min", 0), task.get("budget_max", 0))
|
||
commission = int(amount * COMMISSION_RATE)
|
||
esc = await db.escrows.create(task_id, task["title"], amount, commission)
|
||
await write_audit(db, action="escrow.deposit", resource="escrow", resource_id=esc["id"],
|
||
detail=f"amount={amount}", user=_u, request=request)
|
||
return esc
|
||
|
||
|
||
@router.post("/tasks/{task_id}/release", summary="运营结算(验收通过,划佣金+代付 OPC)")
|
||
async def operator_escrow_release(
|
||
task_id: str,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
_u: dict = Depends(require_roles("operator")),
|
||
):
|
||
"""验收通过:任务须 completed;托管结算 released(平台佣金留存,余额代付 OPC)。"""
|
||
from ...services.settlement_service import SettlementService
|
||
|
||
released = await SettlementService(db).release_escrow(task_id)
|
||
await write_audit(db, action="escrow.release", resource="escrow", resource_id=released["id"],
|
||
detail=f"amount={released.get('amount')}", user=_u, request=request)
|
||
return released
|
||
|
||
|
||
@router.get("/tasks/{task_id}", summary="任务详情")
|
||
async def get_task_detail(
|
||
task_id: str,
|
||
db: Database = Depends(get_db),
|
||
_u: dict = Depends(require_roles("operator")),
|
||
):
|
||
task = await db.tasks.get(task_id)
|
||
if task is None:
|
||
raise HTTPException(status_code=404, detail="任务不存在")
|
||
return task
|
||
|
||
|
||
@router.get("/tasks/{task_id}/bids", summary="任务竞标(报名/投标情况)")
|
||
async def get_task_bids(
|
||
task_id: str,
|
||
db: Database = Depends(get_db),
|
||
_u: dict = Depends(require_roles("operator")),
|
||
):
|
||
return {"items": await db.bids.list_for_task(task_id)}
|
||
|
||
|
||
@router.get("/tasks/{task_id}/claims", summary="任务报名/接单情况(抢单/指派/推荐)")
|
||
async def get_task_claims(
|
||
task_id: str,
|
||
db: Database = Depends(get_db),
|
||
_u: dict = Depends(require_roles("operator")),
|
||
):
|
||
return {"items": await db.task_claims.list_by_task(task_id)}
|
||
|
||
|
||
@router.post("/tasks/{task_id}/bids/{bid_id}/win", summary="开标/定标(选定中标 OPC)")
|
||
async def operator_win_bid(
|
||
task_id: str,
|
||
bid_id: str,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
_u: dict = Depends(require_roles("operator")),
|
||
):
|
||
from ...services.task_service import TaskService
|
||
|
||
bid = await db.bids.set_status(bid_id, "win")
|
||
if bid is None:
|
||
raise HTTPException(status_code=404, detail="竞标不存在")
|
||
# 其余标置为 lost
|
||
for b in await db.bids.list_for_task(task_id):
|
||
if b["id"] != bid_id and b.get("status") == "submitted":
|
||
await db.bids.set_status(b["id"], "lost")
|
||
await db.tasks.set_status(task_id, "in_progress")
|
||
await write_audit(db, action="bid.win", resource="bid", resource_id=bid_id,
|
||
detail=f"task={task_id} winner={bid.get('opc_name')}", user=_u, request=request)
|
||
return bid
|
||
|
||
|
||
@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")),
|
||
):
|
||
fields = req.model_dump(exclude_none=True)
|
||
fields["mode"] = _normalize_mode(fields.get("mode", "grab"))
|
||
if not fields.get("task_code"):
|
||
fields["task_code"] = _gen_task_code()
|
||
task = await db.tasks.create(fields)
|
||
await write_audit(db, action="task.create", resource="task", resource_id=task["id"],
|
||
detail=task["title"], user=actor, request=request)
|
||
return task
|
||
|
||
|
||
def _normalize_mode(mode: str) -> str:
|
||
"""接单方式归一:grab/bid/assign/recommend;旧值 designated→assign、dispatch→recommend。"""
|
||
mapping = {"designated": "assign", "dispatch": "recommend"}
|
||
m = mapping.get(mode, mode)
|
||
return m if m in ("grab", "bid", "assign", "recommend") else "grab"
|
||
|
||
|
||
@router.patch("/tasks/{task_id}", summary="更新任务(增改)")
|
||
async def update_task(
|
||
task_id: str,
|
||
req: TaskUpdateRequest,
|
||
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.update(task_id, req.model_dump(exclude_none=True))
|
||
await write_audit(db, action="task.update", resource="task", resource_id=task_id,
|
||
detail=(req.title or ""), user=actor, request=request)
|
||
return task
|
||
|
||
|
||
@router.get("/task-categories", summary="任务分类字典")
|
||
async def list_task_categories(
|
||
db: Database = Depends(get_db),
|
||
_u: dict = Depends(require_roles("operator")),
|
||
):
|
||
return await db.task_categories.list()
|
||
|
||
|
||
@router.post("/tasks/{task_id}/assign", summary="指派任务给某人")
|
||
async def assign_task(
|
||
task_id: str,
|
||
req: TaskAssignRequest,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
actor: dict = Depends(require_permission("action:task.manage")),
|
||
):
|
||
from ...services.task_service import TaskService
|
||
|
||
updated = await TaskService(db).assign(task_id, req.taker_user_id, actor)
|
||
await write_audit(db, action="task.assign", resource="task", resource_id=task_id,
|
||
detail=req.taker_user_id, user=actor, request=request)
|
||
return updated
|
||
|
||
|
||
@router.post("/tasks/{task_id}/recommend", summary="推荐候选人")
|
||
async def recommend_task(
|
||
task_id: str,
|
||
req: TaskRecommendRequest,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
actor: dict = Depends(require_permission("action:task.manage")),
|
||
):
|
||
from ...services.task_service import TaskService
|
||
|
||
records = await TaskService(db).recommend(task_id, req.candidates, actor)
|
||
await write_audit(db, action="task.recommend", resource="task", resource_id=task_id,
|
||
detail=",".join(req.candidates), user=actor, request=request)
|
||
return {"items": records}
|
||
|
||
|
||
@router.post("/tasks/{task_id}/select", summary="选定推荐人")
|
||
async def select_recommend_task(
|
||
task_id: str,
|
||
req: TaskSelectRecommendRequest,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
actor: dict = Depends(require_permission("action:task.manage")),
|
||
):
|
||
from ...services.task_service import TaskService
|
||
|
||
updated = await TaskService(db).select_recommend(task_id, req.taker_user_id, actor)
|
||
await write_audit(db, action="task.select", resource="task", resource_id=task_id,
|
||
detail=req.taker_user_id, user=actor, request=request)
|
||
return updated
|
||
|
||
|
||
def _gen_task_code() -> str:
|
||
"""生成便于扫码展示的短码:TK-YYYYMMDD-XXXX(基于时间戳短采样)。"""
|
||
import time as _t
|
||
from datetime import datetime as _dt
|
||
|
||
stamp = _dt.now().strftime("%Y%m%d")
|
||
rand = f"{int(_t.time()) % 1000000:06d}"
|
||
return f"TK-{stamp}-{rand}"
|
||
|
||
|
||
@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.publish(task_id)) if req.status == "published" else 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.get("/content/{content_id}", summary="内容详情(编辑回填)")
|
||
async def get_content(
|
||
content_id: str,
|
||
db: Database = Depends(get_db),
|
||
_u: dict = Depends(require_roles("operator")),
|
||
):
|
||
item = await db.content.get(content_id)
|
||
if item is None:
|
||
raise HTTPException(status_code=404, detail="Content not found")
|
||
return item
|
||
|
||
|
||
@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")),
|
||
):
|
||
payload = req.model_dump(exclude_none=True)
|
||
# 发布人:未显式指定时用当前操作者(昵称/头像快照)
|
||
payload.setdefault("publisher_id", actor.get("id", ""))
|
||
if not payload.get("publisher_name"):
|
||
u = await db.users.get_by_id(actor.get("id", ""))
|
||
payload["publisher_name"] = (u or {}).get("nickname") or actor.get("username") or "平台"
|
||
payload["publisher_avatar"] = (u or {}).get("avatar", "") or ""
|
||
item = await db.content.create(payload)
|
||
await write_audit(db, action="content.create", resource="content", resource_id=item["id"],
|
||
detail=item["title"], user=actor, request=request)
|
||
return item
|
||
|
||
|
||
@router.post("/upload", summary="上传媒体(图片/pdf/doc/docx)")
|
||
async def admin_upload(
|
||
file: UploadFile = File(...),
|
||
actor: dict = Depends(require_permission("action:content.manage")),
|
||
):
|
||
"""运营端上传资讯封面 / 图集 / 视频封面等:存 uploads 目录,返回 {ok, url}。"""
|
||
from ...services import media_upload
|
||
url = await media_upload.save_media(file)
|
||
return {"ok": True, "url": url}
|
||
|
||
|
||
@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.post("/content/{content_id}/review", summary="资讯审核(通过/下架)")
|
||
async def review_content(
|
||
content_id: str,
|
||
req: ContentReviewRequest,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
actor: dict = Depends(require_permission("action:content.manage")),
|
||
):
|
||
"""approve → 已发布 + 公开(C 端可见);offline → 已下线。"""
|
||
if req.decision == "approve":
|
||
item = await db.content.approve(content_id)
|
||
action = "content.approve"
|
||
elif req.decision == "offline":
|
||
item = await db.content.set_status(content_id, "offline")
|
||
action = "content.offline"
|
||
else:
|
||
raise HTTPException(status_code=400, detail="无效审核操作")
|
||
if item is None:
|
||
raise HTTPException(status_code=404, detail="Content not found")
|
||
await write_audit(db, action=action, resource="content", resource_id=content_id,
|
||
detail=item.get("title", ""), user=actor, request=request)
|
||
return item
|
||
|
||
|
||
@router.put("/content/{content_id}", summary="编辑内容")
|
||
async def update_content(
|
||
content_id: str,
|
||
req: ContentUpdateRequest,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
actor: dict = Depends(require_permission("action:content.manage")),
|
||
):
|
||
item = await db.content.update(content_id, req.model_dump(exclude_unset=True, exclude_none=True))
|
||
if item is None:
|
||
raise HTTPException(status_code=404, detail="Content not found")
|
||
await write_audit(db, action="content.update", resource="content", resource_id=content_id,
|
||
detail=item.get("title", ""), user=actor, request=request)
|
||
return item
|
||
|
||
|
||
@router.delete("/content/{content_id}", summary="删除内容")
|
||
async def delete_content(
|
||
content_id: str,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
actor: dict = Depends(require_permission("action:content.manage")),
|
||
):
|
||
ok = await db.content.delete(content_id)
|
||
if not ok:
|
||
raise HTTPException(status_code=404, detail="Content not found")
|
||
await write_audit(db, action="content.delete", resource="content", resource_id=content_id,
|
||
detail=content_id, user=actor, request=request)
|
||
return {"ok": True}
|
||
|
||
|
||
# ── 系统配置 ──────────────────────────────────────────────────────────────
|
||
@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"),
|
||
quota_per_unit=data.get("quota_per_unit"),
|
||
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/user-balance", summary="充值/调整引擎用户余额(add_quota)")
|
||
async def compute_user_balance(
|
||
req: ComputeBalanceRequest,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
actor: dict = Depends(require_roles("operator")),
|
||
):
|
||
"""给引擎用户充值/扣减/覆盖余额(配额单位)。
|
||
|
||
value 为配额单位:/api/status 的 quota_per_unit(默认 500000 = $1)。
|
||
mode:add 充值 / subtract 扣减 / override 直接设为该值。
|
||
引擎按用户余额计量,额度耗尽时调用 /v1 返回 403。
|
||
"""
|
||
if req.mode not in ("add", "subtract", "override"):
|
||
raise HTTPException(status_code=400, detail="mode 需为 add/subtract/override")
|
||
if req.mode != "override" and req.value <= 0:
|
||
raise HTTPException(status_code=400, detail="调整量需大于 0")
|
||
try:
|
||
result = await compute_client.adjust_user_quota(req.engine_user_id, req.value, req.mode)
|
||
except compute_client.ComputeError as exc:
|
||
raise HTTPException(status_code=502, detail=f"算力引擎对接失败: {exc}") from exc
|
||
await write_audit(db, action="compute.balance", resource="compute",
|
||
resource_id=str(req.engine_user_id),
|
||
detail=f"{req.mode} {req.value} quota" + (f" ({req.reason})" if req.reason else ""),
|
||
user=actor, request=request)
|
||
return {"ok": result.get("success", True), "engine_user_id": req.engine_user_id,
|
||
"mode": req.mode, "value": req.value, "message": result.get("message", "")}
|
||
|
||
|
||
@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
|
||
|
||
|
||
@router.put("/courses/{course_id}", summary="编辑课程")
|
||
async def update_course(
|
||
course_id: str,
|
||
req: CourseUpdateRequest,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
actor: dict = Depends(require_permission("action:course.manage")),
|
||
):
|
||
course = training.update_course(course_id, req.model_dump(exclude_unset=True, exclude_none=True))
|
||
if course is None:
|
||
raise HTTPException(status_code=404, detail="Course not found")
|
||
await write_audit(db, action="course.update", resource="course", resource_id=course_id,
|
||
detail=course.get("title", ""), user=actor, request=request)
|
||
return course
|
||
|
||
|
||
@router.delete("/courses/{course_id}", summary="删除课程")
|
||
async def delete_course(
|
||
course_id: str,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
actor: dict = Depends(require_permission("action:course.manage")),
|
||
):
|
||
ok = training.delete_course(course_id)
|
||
if not ok:
|
||
raise HTTPException(status_code=404, detail="Course not found")
|
||
await write_audit(db, action="course.delete", resource="course", resource_id=course_id,
|
||
detail=course_id, user=actor, request=request)
|
||
return {"ok": True}
|
||
|
||
|
||
# ── 培训业务 · 活动(桥接 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", summary="创建活动")
|
||
async def create_activity(
|
||
req: ActivityCreateRequest,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
actor: dict = Depends(require_permission("action:activity.manage")),
|
||
):
|
||
activity = training.create_activity(req.model_dump(exclude_none=True))
|
||
await write_audit(db, action="activity.create", resource="activity", resource_id=activity["id"],
|
||
detail=activity["title"], user=actor, request=request)
|
||
return activity
|
||
|
||
|
||
@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
|
||
|
||
|
||
@router.put("/activities/{event_id}", summary="编辑活动")
|
||
async def update_activity(
|
||
event_id: str,
|
||
req: ActivityUpdateRequest,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
actor: dict = Depends(require_permission("action:activity.manage")),
|
||
):
|
||
activity = training.update_activity(event_id, req.model_dump(exclude_unset=True, exclude_none=True))
|
||
if activity is None:
|
||
raise HTTPException(status_code=404, detail="Activity not found")
|
||
await write_audit(db, action="activity.update", resource="activity", resource_id=event_id,
|
||
detail=activity.get("title", ""), user=actor, request=request)
|
||
return activity
|
||
|
||
|
||
@router.delete("/activities/{event_id}", summary="删除活动")
|
||
async def delete_activity(
|
||
event_id: str,
|
||
request: Request,
|
||
db: Database = Depends(get_db),
|
||
actor: dict = Depends(require_permission("action:activity.manage")),
|
||
):
|
||
ok = training.delete_activity(event_id)
|
||
if not ok:
|
||
raise HTTPException(status_code=404, detail="Activity not found")
|
||
await write_audit(db, action="activity.delete", resource="activity", resource_id=event_id,
|
||
detail=event_id, user=actor, request=request)
|
||
return {"ok": True}
|
||
|
||
|
||
# ── 培训业务 · 报名(桥接 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
|
||
|
||
|
||
# ── 培训业务 · 调研 / 政策 / 流程日志(落库结果只读)─────────────────────────
|
||
@router.get("/surveys", summary="调研结果列表")
|
||
async def list_surveys(
|
||
_u: dict = Depends(require_roles("operator")),
|
||
):
|
||
return training.list_surveys()
|
||
|
||
|
||
@router.get("/policies", summary="政策测评结果列表")
|
||
async def list_policies(
|
||
_u: dict = Depends(require_roles("operator")),
|
||
):
|
||
return training.list_policies()
|
||
|
||
|
||
@router.get("/plans", summary="流程启动日志列表")
|
||
async def list_plans(
|
||
_u: dict = Depends(require_roles("operator")),
|
||
):
|
||
return training.list_plans()
|