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

1213 lines
49 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 -*-
"""运营端业务管理端点:任务 / 服务商 / 内容 / 系统配置 / 数据统计。
全部要求 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
transaction_id: str = "" # 微信支付单号(服务商代子商户收单支付成功后回填)
payer_sub_mchid: str = "" # 出资特约商户号(收单子商户,即分账出资方)
# ── OPC 收款认证(桌面端开户认证管理:进件/签约/账户验证/分账授权)────────────────
@router.get("/pay/bindings", summary="OPC 收款绑定列表(进件状态/分账授权总览)")
async def operator_pay_bindings(
bind_type: str | None = None,
status: str | None = None,
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("operator")),
):
"""管理端:全部 OPC 收款绑定(个人 openid / 特约商户进件),带用户名、进件状态、签约链接、分账授权。"""
from ...pay.repository import PaymentBindingRepository
from ...infrastructure.repositories import UserRepository
items = await PaymentBindingRepository(db.session).list_all(
bind_type=bind_type, status=status)
users = {u["id"]: u for u in await UserRepository(db.session).list()}
enriched = []
for b in items:
u = users.get(b["user_id"], {})
detail = {}
try:
import json as _json
detail = _json.loads(b.get("audit_detail_json") or "{}")
except Exception:
detail = {}
enriched.append({
**b,
"user_name": u.get("nickname") or u.get("username") or "",
"user_phone": u.get("phone") or u.get("username") or "",
"audit_detail": detail,
})
return {"items": enriched}
@router.post("/pay/bindings/{binding_id}/refresh", summary="代查进件状态(同步微信侧,超管签约引导)")
async def operator_refresh_applyment(
binding_id: str,
request: Request,
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("operator")),
):
"""服务商替 OPC 查询进件状态并回填:sign_url(超管扫码验证+签约)、sub_mchid、账户验证、
驳回原因;FINISHED 自动添加分账接收方;查询前自动查一次分账授权。"""
from ...pay import profitsharing
from ...pay.repository import PaymentBindingRepository
repo = PaymentBindingRepository(db.session)
b = await repo.get(binding_id)
if b is None:
raise HTTPException(status_code=404, detail="绑定不存在")
applyment_id = b.get("applyment_id") or ""
if not applyment_id:
raise HTTPException(status_code=400, detail="该绑定未发起进件")
data = await profitsharing.query_applyment(applyment_id=applyment_id)
if data is None:
raise HTTPException(status_code=503, detail="查询进件状态失败,请稍后重试")
state = data.get("applyment_state") or ""
sub_mchid = str(data.get("sub_mchid") or "")
sign_url = str(data.get("sign_url") or "")
fields = {"applyment_state": state}
if sign_url:
fields["sign_url"] = sign_url
if sub_mchid:
fields["sub_mchid"] = sub_mchid
if data.get("account_validation"):
fields["account_validation_json"] = json.dumps(data["account_validation"], ensure_ascii=False)
audit = {}
if data.get("audit_detail"):
audit["audit_detail"] = data["audit_detail"]
if data.get("applyment_state_msg"):
audit["applyment_state_msg"] = data["applyment_state_msg"]
if audit:
fields["audit_detail_json"] = json.dumps(audit, ensure_ascii=False)
if state == "FINISHED" and sub_mchid:
try:
await profitsharing.add_receiver(
account_type="MERCHANT_ID", account=sub_mchid, sub_mchid=sub_mchid)
fields.update({"status": "active", "detail": "进件通过,已添加分账接收方"})
except RuntimeError as exc:
fields.update({"status": "failed", "detail": f"添加接收方失败:{str(exc)[:200]}"})
# 顺带查询分账授权(子商户是否已开通分账)
cfg = await profitsharing.query_split_config(sub_mchid=sub_mchid)
if cfg and cfg.get("max_ratio") is not None:
fields["split_allowed"] = "OPEN" if int(cfg["max_ratio"]) > 0 else "CLOSED"
fields["split_max_ratio"] = int(cfg["max_ratio"])
elif state == "REJECTED":
reason = ""
ad = data.get("audit_detail")
if isinstance(ad, dict) and isinstance(ad.get("reject_reason"), list):
reason = "".join(str(x) for x in ad["reject_reason"])
fields.update({"status": "rejected", "detail": reason or data.get("applyment_state_msg") or "进件被驳回"})
elif state in profitsharing.APPLYMENT_ACTION_STATES:
fields["detail"] = data.get("applyment_state_msg") or \
profitsharing.APPLYMENT_STATE.get(state, state)
b = await repo.update(binding_id, fields)
await write_audit(db, action="pay.binding.refresh", resource="payment_binding",
resource_id=binding_id,
detail=f"state={state} sub_mchid={sub_mchid}", user=_u, request=request)
return {"binding": b, "applyment_state": state, "sign_url": sign_url or b.get("sign_url"),
"sub_mchid": sub_mchid or b.get("sub_mchid")}
@router.get("/pay/bindings/{binding_id}/split-config", summary="查询 OPC 子商户分账授权/最大比例")
async def operator_split_config(
binding_id: str,
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("operator")),
):
from ...pay import profitsharing
from ...pay.repository import PaymentBindingRepository
repo = PaymentBindingRepository(db.session)
b = await repo.get(binding_id)
if b is None:
raise HTTPException(status_code=404, detail="绑定不存在")
if not b.get("sub_mchid"):
raise HTTPException(status_code=400, detail="该绑定无子商户号,无法查询分账授权")
cfg = await profitsharing.query_split_config(sub_mchid=b["sub_mchid"])
if cfg is None:
raise HTTPException(status_code=503, detail="子商户未开通分账或查询失败,请确认子商户在商户平台开通分账并授权")
max_ratio = int(cfg.get("max_ratio") or 0)
b = await repo.update(binding_id, {
"split_allowed": "OPEN" if max_ratio > 0 else "CLOSED",
"split_max_ratio": max_ratio,
})
return {"sub_mchid": b["sub_mchid"], "split_allowed": b["split_allowed"],
"split_max_ratio": max_ratio, "split_max_ratio_percent": round(max_ratio / 100.0, 1),
"binding": b}
@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")),
):
"""发标方向运营平台托管账户建档/充值;已托管返回现有。
传入 transaction_id + payer_sub_mchid 时(服务商代子商户收单支付成功),
记录微信支付单号并置 frozen(资金已在微信侧冻结,可分账)。
"""
from ...services.settlement_service import COMMISSION_RATE, SettlementService
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:
if req.transaction_id and not e.get("transaction_id"):
e = await SettlementService(db).release_escrow(
task_id, transaction_id=req.transaction_id,
payer_sub_mchid=req.payer_sub_mchid)
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)
if req.transaction_id:
esc = await SettlementService(db).release_escrow(
task_id, transaction_id=req.transaction_id,
payer_sub_mchid=req.payer_sub_mchid)
await write_audit(db, action="escrow.deposit", resource="escrow", resource_id=esc["id"],
detail=f"amount={amount} tx={req.transaction_id}", user=_u, request=request)
return esc
@router.post("/tasks/{task_id}/release", summary="运营结算(验收通过,微信分账+解冻剩余 或 台账降级)")
async def operator_escrow_release(
task_id: str,
request: Request,
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("operator")),
):
"""验收通过:任务须 completed;优先走微信服务商分账完整生命周期
(查询剩余待分→分账平台佣金5%→结果确认→解冻剩余95%给出资方),
未接微信收单/无分账配置时降级台账结算(released)。"""
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')} share_status={released.get('share_status')}",
user=_u, request=request)
return released
@router.get("/tasks/{task_id}/escrow/status", summary="任务托管/分账生命周期状态")
async def operator_escrow_status(
task_id: str,
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("operator")),
):
"""查询任务托管 + 分账生命周期(冻结→分账→解冻剩余),分账处理中时主动查询微信侧最新结果。"""
from ...services.settlement_service import SettlementService
return await SettlementService(db).query_split_lifecycle(task_id)
@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
# ── 内容管理 ──────────────────────────────────────────────────────────────
def _resolve_content_media(items: list[dict]) -> list[dict]:
"""运营端出口:封面/视频/发布人头像等对象路径统一生成 CDN 直链。"""
from ...infrastructure.oss import resolve_url
for it in items:
for k in ("cover", "cover_small", "video", "video_cover", "publisher_avatar"):
it[k] = resolve_url(it.get(k, ""))
it["images"] = [resolve_url(u) for u in (it.get("images") or [])]
return items
@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 _resolve_content_media(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 _resolve_content_media([item])[0]
@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"):
payload["publisher_name"] = "官方"
payload["publisher_avatar"] = ""
elif payload.get("publisher_avatar"):
from ...infrastructure.oss import to_object_path
payload["publisher_avatar"] = to_object_path(payload["publisher_avatar"])
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")),
):
"""运营端上传资讯封面 / 图集 / 视频封面等:统一走 OSS(news/ 业务目录),返回 {ok, url}。"""
from ...services import media_upload
url = await media_upload.save_media(file, dir="news")
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 为微元(人民币实际金额,1 元 = 1_000_000compute-service meter.py 口径)。
modeadd 充值 / 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
# 充值后即时回写平台侧余额镜像,避免与引擎实际余额漂移(best-effort)。
await compute_client.sync_user_mirror(db, req.engine_user_id)
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.get("/compute/platform-users", summary="平台用户 × 引擎真实余额(充值以引擎为准)")
async def compute_platform_users(
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("operator")),
):
"""按「平台用户」聚合展示算力账户,余额取**引擎真实值**compute.db users.quota,扣费真源)。
- 映射键:用户名字符串(platform.username == engine.username)。
- 引擎侧字段(engine_id / quota / used_quota / discount / group)逐页全量拉取后 join
引擎缺失的用户标记 provisioned=False(余额回退镜像值并置 0)。
- 镜像(compute_quota)仅作参考列返回,**不作为充值/展示真源**。
"""
users = await db.users.list()
engine: dict[str, dict] = {}
page = 1
while page <= 20: # 上限 20 页 × 100 = 2000 引擎用户,足够且防失控
try:
items = await compute_client.list_engine_users(page=page, page_size=100)
except compute_client.ComputeError as exc:
raise HTTPException(status_code=502, detail=f"算力引擎对接失败: {exc}") from exc
for it in items:
uname = str(it.get("username") or "")
if uname:
engine[uname] = it
if len(items) < 100:
break
page += 1
items_out = []
for u in users:
uname = u.get("username") or ""
e = engine.get(uname) or {}
items_out.append({
"id": u["id"],
"username": uname,
"nickname": u.get("nickname") or "",
"role": u.get("role") or "",
"provisioned": bool(e),
"engine_id": e.get("id"),
# 真实余额(引擎为准);引擎缺失时回退镜像(0)
"quota": int(e.get("quota") if e else (u.get("compute_quota") or 0)) or 0,
"used_quota": int(e.get("used_quota") if e else (u.get("compute_used_quota") or 0)) or 0,
"discount": int(e.get("discount") or 100) if e else 100,
"group": e.get("group") or "",
"engine_display_name": e.get("display_name") or "",
"request_count": int(e.get("request_count") or 0) if e else 0,
"last_login_at": e.get("last_login_at"),
"mirror_quota": int(u.get("compute_quota") or 0), # 参考列
})
return {"items": items_out, "total": len(items_out), "engine_total": len(engine)}
@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()
# ── 汇总统计 & 题目定义 ────────────────────────────────────────────────────
@router.get("/surveys/questions", summary="调研题目定义")
async def get_survey_questions(
_u: dict = Depends(require_roles("operator")),
):
return training.survey_questions()
@router.get("/surveys/summary", summary="调研汇总统计")
async def get_survey_summary(
_u: dict = Depends(require_roles("operator")),
):
return training.survey_summary()
@router.get("/tests/summary", summary="测评汇总统计")
async def get_test_summary(
_u: dict = Depends(require_roles("operator")),
):
return training.test_summary()
@router.get("/bookings/summary", summary="报名汇总统计")
async def get_booking_summary(
_u: dict = Depends(require_roles("operator")),
):
return training.booking_summary()
@router.get("/policies/summary", summary="政策测评汇总统计")
async def get_policy_summary(
_u: dict = Depends(require_roles("operator")),
):
return training.policy_summary()
@router.get("/plans/summary", summary="流程启动汇总统计")
async def get_plan_summary(
_u: dict = Depends(require_roles("operator")),
):
return training.plan_summary()
# ── 各模块详情 / 删除(统一模式)──────────────────────────────────────────
@router.get("/tests/{test_id}", summary="测评详情")
async def get_test(
test_id: str,
_u: dict = Depends(require_roles("operator")),
):
item = training.get_test(test_id)
if item is None:
raise HTTPException(status_code=404, detail="Test not found")
return item
@router.delete("/tests/{test_id}", summary="删除测评记录")
async def delete_test(
test_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:test.manage")),
):
if not training.delete_test(test_id):
raise HTTPException(status_code=404, detail="Test not found")
await write_audit(db, action="test.delete", resource="test", resource_id=test_id,
detail=test_id, user=actor, request=request)
return {"ok": True}
@router.get("/bookings/{booking_id}", summary="报名详情")
async def get_booking(
booking_id: str,
_u: dict = Depends(require_roles("operator")),
):
item = training.get_booking(booking_id)
if item is None:
raise HTTPException(status_code=404, detail="Booking not found")
return item
@router.delete("/bookings/{booking_id}", summary="删除报名")
async def delete_booking(
booking_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:booking.manage")),
):
if not training.delete_booking(booking_id):
raise HTTPException(status_code=404, detail="Booking not found")
await write_audit(db, action="booking.delete", resource="booking", resource_id=booking_id,
detail=booking_id, user=actor, request=request)
return {"ok": True}
@router.get("/surveys/{survey_id}", summary="调研详情")
async def get_survey(
survey_id: str,
_u: dict = Depends(require_roles("operator")),
):
item = training.get_survey(survey_id)
if item is None:
raise HTTPException(status_code=404, detail="Survey not found")
return item
@router.delete("/surveys/{survey_id}", summary="删除调研记录")
async def delete_survey(
survey_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:survey.manage")),
):
if not training.delete_survey(survey_id):
raise HTTPException(status_code=404, detail="Survey not found")
await write_audit(db, action="survey.delete", resource="survey", resource_id=survey_id,
detail=survey_id, user=actor, request=request)
return {"ok": True}
@router.get("/policies/{policy_id}", summary="政策测评详情")
async def get_policy(
policy_id: str,
_u: dict = Depends(require_roles("operator")),
):
item = training.get_policy(policy_id)
if item is None:
raise HTTPException(status_code=404, detail="Policy not found")
return item
@router.delete("/policies/{policy_id}", summary="删除政策测评记录")
async def delete_policy(
policy_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:policy.manage")),
):
if not training.delete_policy(policy_id):
raise HTTPException(status_code=404, detail="Policy not found")
await write_audit(db, action="policy.delete", resource="policy", resource_id=policy_id,
detail=policy_id, user=actor, request=request)
return {"ok": True}
@router.get("/plans/{plan_id}", summary="流程启动详情")
async def get_plan(
plan_id: str,
_u: dict = Depends(require_roles("operator")),
):
item = training.get_plan(plan_id)
if item is None:
raise HTTPException(status_code=404, detail="Plan not found")
return item
@router.delete("/plans/{plan_id}", summary="删除流程启动记录")
async def delete_plan(
plan_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_permission("action:plan.manage")),
):
if not training.delete_plan(plan_id):
raise HTTPException(status_code=404, detail="Plan not found")
await write_audit(db, action="plan.delete", resource="plan", resource_id=plan_id,
detail=plan_id, user=actor, request=request)
return {"ok": True}