Files
server-core/app/api/routers/rbac_opc.py
T
Pine 873a70a095 feat(content): 资讯增删改 + C端资讯中心(类别)端点
- ContentRepository:+update(部分字段)/get/delete
- schemas.operator:+ContentUpdateRequest(全可选)
- rbac_operator:+PUT/DELETE /admin/content/{id}(action:content.manage+审计)
- rbac_opc:+GET /opc/content?type=policy|news|skill|dynamic(仅published)、/opc/content/{id}(详情)
- 类别字典 policy/news/skill/dynamic

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-27 14:05:38 +08:00

500 lines
19 KiB
Python

# -*- coding: utf-8 -*-
"""OPC 超级个体端端点(工作台聚合数据)。
所有数据均由服务端(SQLite 演示库)提供,前端不落地任何演示数据。
当前角色要求 opc_member;后续各页面端点在此扩展。
"""
from __future__ import annotations
from fastapi import APIRouter, Depends, HTTPException, Request
from pydantic import BaseModel
from ..dependencies import get_db
from ..schemas.opc import BidRequest, ProfileUpdate, FinanceRecordCreate, TaskClaimRequest, TaskAssignRequest, TaskRecommendRequest, TaskSelectRecommendRequest
from ...rbac import require_roles, write_audit
from ...infrastructure.repositories import Database, new_id, utcnow_iso
from ...infrastructure.models import FinanceRecord
from ... import config
from ...services import compute_client, compute_catalog
router = APIRouter(prefix="/opc", tags=["opc"])
@router.get("/dashboard", summary="OPC 工作台总览")
async def opc_dashboard(
db: Database = Depends(get_db),
user: dict = Depends(require_roles("opc_member")),
):
"""返回当前 OPC 用户的工作台聚合数据(统计/进行中任务/智能体建议/月度收入/最新消息)。"""
from ...services.dashboard_service import DashboardService
return await DashboardService(db).opc_dashboard(user["id"])
# ── OPC 子页面(全部由服务端提供标准数据)─────────────────────────────────
@router.get("/agent", summary="OPC 智能体助手")
async def opc_agent(
db: Database = Depends(get_db),
user: dict = Depends(require_roles("opc_member")),
):
agents = await db.agents.get_by_user(user["id"], port=user.get("port") or "opc")
default = next((a for a in agents if a.get("id", "").startswith("pine_agents_official_001")), None)
return {
"agent": {
"name": (default or {}).get("name", "小云"),
"desc": (default or {}).get("description", "您的智能体助手,可处理创业、政策、财税等事务"),
"port": user.get("port") or "opc",
},
"quick_actions": [
{"title": "查询政策", "desc": "为您匹配最新创业政策"},
{"title": "发布能力", "desc": "快速发布您的服务能力"},
{"title": "财务助手", "desc": "查看收入与支出"},
],
"chat_path": "/chat",
}
@router.get("/tasks", summary="任务广场(已发布平台任务)")
async def opc_task_square(
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("opc_member")),
):
return {"items": await db.tasks.list(status="published")}
@router.get("/my-tasks", summary="我的任务看板")
async def opc_my_tasks(
db: Database = Depends(get_db),
user: dict = Depends(require_roles("opc_member")),
):
return {"items": await db.opc_tasks.list_by_user(user["id"])}
@router.get("/services", summary="服务市场(在营服务商)")
async def opc_services(
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("opc_member")),
):
return {"items": await db.providers.list(status="active")}
@router.get("/policy", summary="政策列表(已发布政策)")
async def opc_policy(
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("opc_member")),
):
return {"items": await db.content.list(ctype="policy", status="published")}
@router.get("/content", summary="资讯中心(按类别,已发布)")
async def opc_content(
type: str = "news",
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("opc_member")),
):
"""C 端资讯:type ∈ policy(政策)/news(资讯)/skill(技能)/dynamic(动态),仅 published。"""
return {"items": await db.content.list(ctype=type, status="published")}
@router.get("/content/{content_id}", summary="资讯详情")
async def opc_content_detail(
content_id: str,
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("opc_member")),
):
item = await db.content.get(content_id)
if item is None or item.get("status") != "published":
raise HTTPException(status_code=404, detail="内容不存在")
return item
@router.get("/finance", summary="财务流水")
async def opc_finance(
db: Database = Depends(get_db),
user: dict = Depends(require_roles("opc_member")),
):
uid = user["id"]
return {
"records": await db.finance.list_by_user(uid),
"total_income": await db.finance.total_income(uid),
}
@router.get("/messages", summary="站内消息")
async def opc_messages(
db: Database = Depends(get_db),
user: dict = Depends(require_roles("opc_member")),
):
return {"items": await db.messages.list_by_user(user["id"])}
@router.get("/profile", summary="个人资料")
async def opc_profile(
db: Database = Depends(get_db),
user: dict = Depends(require_roles("opc_member")),
):
profile = await db.users.to_profile(user)
profile["opc_profile"] = await db.opc_profiles.get(user["id"])
return profile
@router.post("/tasks/{task_id}/grab", summary="OPC 抢单")
async def opc_grab(
task_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("opc_member")),
):
from ...services.task_service import TaskService
updated = await TaskService(db).grab(task_id, actor)
await write_audit(db, action="task.grab", resource="task", resource_id=task_id,
detail=actor.get("username"), user=actor, request=request)
return updated
@router.post("/tasks/grab-by-code", summary="扫码/按短码接单")
async def opc_grab_by_code(
req: TaskClaimRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("opc_member")),
):
from ...services.task_service import TaskService
updated = await TaskService(db).claim(
(req.task_code or "").strip(), actor, source="scan",
)
await write_audit(db, action="task.claim", resource="task",
resource_id=updated["id"], detail=actor.get("username"),
user=actor, request=request)
return updated
@router.post("/tasks/{task_id}/doing", summary="开始做单")
async def opc_task_doing(
task_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("opc_member")),
):
from ...services.task_service import TaskService
updated = await TaskService(db).start_doing(task_id, actor)
await write_audit(db, action="task.doing", resource="task", resource_id=task_id,
user=actor, request=request)
return updated
@router.post("/tasks/{task_id}/complete", summary="完成做单")
async def opc_task_complete(
task_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("opc_member")),
):
from ...services.task_service import TaskService
updated = await TaskService(db).complete(task_id, actor)
await write_audit(db, action="task.complete", resource="task", resource_id=task_id,
user=actor, request=request)
return updated
@router.post("/tasks/{task_id}/assign", summary="指派任务给某人")
async def opc_task_assign(
task_id: str,
req: TaskAssignRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("opc_member")),
):
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 opc_task_recommend(
task_id: str,
req: TaskRecommendRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("opc_member")),
):
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 opc_task_select_recommend(
task_id: str,
req: TaskSelectRecommendRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("opc_member")),
):
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
@router.post("/tasks/{task_id}/bid", summary="OPC 投标")
async def opc_bid(
task_id: str,
req: BidRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("opc_member")),
):
from ...services.task_service import TaskService
bid = await TaskService(db).bid(task_id, actor, req.quote, req.plan)
await write_audit(db, action="task.bid", resource="bid", resource_id=bid["id"],
user=actor, request=request)
return bid
@router.post("/tasks/{task_id}/deliver", summary="OPC 交付")
async def opc_deliver(
task_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("opc_member")),
):
from ...services.task_service import TaskService
updated = await TaskService(db).deliver(task_id, actor)
await write_audit(db, action="task.deliver", resource="task", resource_id=task_id,
user=actor, request=request)
return updated
@router.put("/profile", summary="OPC 基础资料编辑")
async def opc_update_profile(
req: ProfileUpdate,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("opc_member")),
):
profile_updates = {k: v for k, v in req.model_dump(exclude_none=True).items()
if k in ("nickname", "account", "company", "room")}
if profile_updates:
await db.users.update_profile(actor["id"], profile_updates)
if req.credit_score is not None:
await db.opc_profiles.upsert(actor["id"], req.credit_score)
await write_audit(db, action="profile.update", resource="user", resource_id=actor["id"],
user=actor, request=request)
fresh_user = await db.users.get_by_id(actor["id"])
return await opc_profile(db, fresh_user)
@router.post("/finance/record", summary="OPC 记一笔收支")
async def opc_add_finance(
req: FinanceRecordCreate,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("opc_member")),
):
from ...services.finance_service import FinanceService
record = await FinanceService(db).add_record(
user_id=actor["id"], category=req.category, amount=req.amount,
date=req.date or utcnow_iso()[:10], note=req.note,
)
await write_audit(db, action="finance.add", resource="finance", resource_id=record["id"],
user=actor, request=request)
return record
@router.get("/tax", summary="OPC 税务申报")
async def opc_tax(
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("opc_member")),
):
return await db.portal_pages.get("opc", "tax") or {"items": []}
@router.get("/affairs", summary="OPC 工商/政务办事")
async def opc_affairs(
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("opc_member")),
):
return await db.portal_pages.get("opc", "affairs") or {"items": []}
# ── 算力中心(用户端,供应商自动对接 PineAgents,令牌自服务)─────────────────
@router.get("/compute/base", summary="算力基础 API 地址")
async def opc_compute_base(
request: Request,
user: dict = Depends(require_roles("opc_member")),
):
"""用户接入地址:经 server-core /v1 中继(OpenAI 兼容)。优先用公网配置,回落请求 host。"""
base = config.COMPUTE_PUBLIC_BASE or str(request.base_url).rstrip("/")
base = base.rstrip("/")
return {
"base_url": f"{base}/v1",
"relay_path": "/v1/chat/completions",
"models_path": "/v1/models",
"username": user.get("username"),
"hint": "将 base_url 与下方令牌填入 OpenAI 兼容客户端(如 base_url/api_key)即可调用本平台模型",
}
@router.get("/compute/models", summary="可用模型(算力中心,admin 新增)")
async def opc_compute_models(_u: dict = Depends(require_roles("opc_member"))):
return {"items": await compute_catalog.models()}
@router.get("/compute/prices", summary="模型价目(admin 模型)")
async def opc_compute_prices(_u: dict = Depends(require_roles("opc_member"))):
return {"items": await compute_catalog.prices()}
@router.get("/compute/usage", summary="我的用量(本月按模型)")
async def opc_compute_usage(user: dict = Depends(require_roles("opc_member"))):
try:
data = await compute_client.user_usage(user.get("username"))
except compute_client.ComputeError as exc:
raise HTTPException(status_code=502, detail=f"算力引擎对接失败: {exc}") from exc
return data
@router.get("/compute/balance", summary="我的算力余额")
async def opc_compute_balance(user: dict = Depends(require_roles("opc_member"))):
try:
data = await compute_client.user_balance(user.get("username"))
except compute_client.ComputeError as exc:
raise HTTPException(status_code=502, detail=f"算力引擎对接失败: {exc}") from exc
return data
@router.get("/compute/tokens", summary="我的算力令牌")
async def opc_compute_tokens(user: dict = Depends(require_roles("opc_member"))):
"""列出当前用户的引擎消费令牌。"""
try:
data = await compute_client.list_user_tokens(user.get("username"))
except compute_client.ComputeError as exc:
raise HTTPException(status_code=502, detail=f"算力引擎对接失败: {exc}") from exc
items = data if isinstance(data, list) else data.get("items", [])
return {"items": items}
@router.post("/compute/tokens", summary="创建算力令牌")
async def opc_compute_tokens_create(user: dict = Depends(require_roles("opc_member"))):
"""为当前用户签发一枚消费令牌(供应商自动对接 PineAgents)。"""
try:
token = await compute_client.issue_user_token(user.get("username"), name="OPC 端令牌")
except compute_client.ComputeError as exc:
raise HTTPException(status_code=502, detail=f"算力引擎对接失败: {exc}") from exc
return token
@router.delete("/compute/tokens/{token_id}", summary="删除算力令牌")
async def opc_compute_tokens_delete(token_id: int, user: dict = Depends(require_roles("opc_member"))):
try:
await compute_client.delete_token(token_id)
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"])}