a60f0ed180
- 迁移0017:content_items.read_count + content_comments 表
- ContentRepository:_to_dict 加 read_count;incr_read_count/content_comments/add_content_comment
- rbac_opc:GET /opc/content/{id} 自增阅读数;GET/POST /opc/content/{id}/comments(发布需 opc_member,带昵称/头像)
Co-Authored-By: Claude <noreply@anthropic.com>
535 lines
21 KiB
Python
535 lines
21 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),
|
||
):
|
||
"""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),
|
||
):
|
||
item = await db.content.get(content_id)
|
||
if item is None or item.get("status") != "published":
|
||
raise HTTPException(status_code=404, detail="内容不存在")
|
||
# 阅读数自增(每次详情浏览 +1)
|
||
if item.get("status") == "published":
|
||
await db.content.incr_read_count(content_id)
|
||
item["read_count"] = (item.get("read_count") or 0) + 1
|
||
return item
|
||
|
||
|
||
@router.get("/content/{content_id}/comments", summary="资讯留言列表")
|
||
async def opc_content_comments(
|
||
content_id: str,
|
||
db: Database = Depends(get_db),
|
||
):
|
||
return {"items": await db.content.content_comments(content_id)}
|
||
|
||
|
||
class ContentCommentBody(BaseModel):
|
||
content: str
|
||
|
||
|
||
@router.post("/content/{content_id}/comments", summary="发表留言")
|
||
async def opc_content_add_comment(
|
||
content_id: str,
|
||
body: ContentCommentBody,
|
||
db: Database = Depends(get_db),
|
||
user: dict = Depends(require_roles("opc_member")),
|
||
):
|
||
if not body.content.strip():
|
||
raise HTTPException(status_code=400, detail="留言不能为空")
|
||
item = await db.content.get(content_id)
|
||
if item is None or item.get("status") != "published":
|
||
raise HTTPException(status_code=404, detail="内容不存在")
|
||
# 取当前用户资料(昵称/头像)用于留言展示
|
||
me = await db.users.get_by_id(user["id"]) or {}
|
||
actor = {"id": user["id"], "username": user.get("username", ""),
|
||
"nickname": (me.get("nickname") or user.get("username", "")),
|
||
"avatar": me.get("avatar", "")}
|
||
comment = await db.content.add_content_comment(content_id, actor, body.content.strip())
|
||
return {"ok": True, "comment": comment}
|
||
|
||
|
||
@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"])}
|