Files
server-core/app/api/routers/rbac_ecosystem.py
T
Pine 980a2db6d9 refactor: 平台应用异步四层架构(接口/业务/领域/基础设施)
- 基础设施层 async:SQLAlchemy 异步引擎/会话、33 Repository async 化、models/security/seed 迁入 infrastructure、新增 cache.py(redis.asyncio) 与 oss.py(aioboto3)
- 接口层:routers 迁 api/routers 并全 async,dependencies 迁 api/dependencies(get_db/get_current_user async)
- 依赖:sqlalchemy[asyncio]/aiosqlite/asyncmy/redis/aioboto3;config 异步 URL + Redis/OSS 配置
- 删除废弃:旧同步 db/dependencies/repositories/storage
- 验证:平台 19 路由 + 培训 48 路由全注册;/health /auth/login /auth/me /admin/tasks /notifications 等接口 async 可用
2026-08-23 23:52:58 +08:00

217 lines
9.2 KiB
Python

# -*- coding: utf-8 -*-
"""横切能力端点:通知/信用/合同/结算/争议/评价/投资撮合/路演参与。
"""
from __future__ import annotations
from fastapi import APIRouter, Depends, HTTPException, Request
from pydantic import BaseModel
from ..dependencies import get_db
from ...rbac import require_permission, require_roles, write_audit
from ...infrastructure.repositories import Database
router = APIRouter(tags=["ecosystem"])
_ANY = ("opc_member", "carrier", "enterprise", "provider", "government", "operator", "investor")
# ── 通知中心 ──────────────────────────────────────────────────────────────
@router.get("/notifications", summary="我的通知")
async def list_notifications(
db: Database = Depends(get_db),
user: dict = Depends(require_roles(*_ANY)),
):
return {"items": await db.notifications.list_for(user["id"]), "unread": await db.notifications.unread(user["id"])}
@router.post("/notifications/read-all", summary="全部已读")
async def read_all_notifications(
db: Database = Depends(get_db),
user: dict = Depends(require_roles(*_ANY)),
):
return {"marked": await db.notifications.mark_read(user["id"])}
# ── 信用/评价 ──────────────────────────────────────────────────────────────
@router.get("/me/credit", summary="我的信用")
async def my_credit(
db: Database = Depends(get_db),
user: dict = Depends(require_roles(*_ANY)),
):
profile = await db.opc_profiles.get(user["id"])
return {"credit_score": (profile or {}).get("credit_score", 80),
"avg_rating": await db.ratings.avg_for(user["id"]),
"rating_count": len(await db.ratings.list_for(user["id"])) if hasattr(db.ratings, "list_for") else 0}
class RatingRequest(BaseModel):
score: int = 5
comment: str = ""
to_id: str = ""
@router.post("/tasks/{task_id}/rate", summary="任务互评")
async def rate_task(
task_id: str,
req: RatingRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("enterprise", "opc_member")),
):
if not req.to_id:
raise HTTPException(status_code=400, detail="缺少被评对象")
rating = await db.ratings.create(task_id, actor["id"], req.to_id, req.score, req.comment)
await write_audit(db, action="task.rate", resource="task", resource_id=task_id,
detail=f"{req.score}", user=actor, request=request)
return rating
# ── 合同 ──────────────────────────────────────────────────────────────────
@router.get("/tasks/{task_id}/contract", summary="查看合同")
async def get_contract(
task_id: str,
db: Database = Depends(get_db),
_u: dict = Depends(require_roles("enterprise", "opc_member")),
):
contract = await db.contracts.get_for_task(task_id)
if contract is None:
raise HTTPException(status_code=404, detail="Contract not found")
return contract
class SignRequest(BaseModel):
opc_id: str = ""
@router.post("/tasks/{task_id}/sign", summary="签订电子合同")
async def sign_contract(
task_id: str,
req: SignRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("enterprise", "opc_member")),
):
task = await db.tasks.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="Task not found")
if await db.contracts.get_for_task(task_id) is not None:
return await db.contracts.get_for_task(task_id)
opc_id = req.opc_id or (actor["id"] if actor.get("role") == "opc_member" else "")
contract = await db.contracts.create(task_id, task["title"], actor["id"], opc_id)
await write_audit(db, action="contract.sign", resource="contract", resource_id=contract["id"],
user=actor, request=request)
return contract
# ── 结算托管 ──────────────────────────────────────────────────────────────
@router.post("/enterprise/tasks/{task_id}/release", summary="验收通过并结算")
async def release_escrow(
task_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("enterprise")),
):
task = await db.tasks.get(task_id)
if task is None or task["status"] != "completed":
raise HTTPException(status_code=400, detail="任务未完成,不能结算")
escrow = await db.escrows.list()
esc = next((e for e in escrow if e["task_id"] == task_id), None)
if esc is None:
amount = max(task.get("budget_min", 0), task.get("budget_max", 0))
commission = int(amount * 0.05)
esc = await db.escrows.create(task_id, task["title"], amount, commission)
released = await db.escrows.set_status(esc["id"], "released")
await write_audit(db, action="escrow.release", resource="escrow", resource_id=esc["id"],
detail=f"amount={esc['amount']}", user=actor, request=request)
return released
@router.get("/operator/settlements", summary="运营结算台账")
async def operator_settlements(
status: str | None = None,
db: Database = Depends(get_db),
_u: dict = Depends(require_permission("action:settlement.manage")),
):
return {"items": await db.escrows.list(status=status)}
# ── 争议 ──────────────────────────────────────────────────────────────────
class DisputeRequest(BaseModel):
reason: str = ""
@router.post("/tasks/{task_id}/dispute", summary="发起争议")
async def create_dispute(
task_id: str,
req: DisputeRequest,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("enterprise", "opc_member")),
):
task = await db.tasks.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="Task not found")
dispute = await db.disputes.create(task_id, task["title"], actor["id"], req.reason)
await db.tasks.set_status(task_id, "disputed")
await write_audit(db, action="dispute.open", resource="dispute", resource_id=dispute["id"],
user=actor, request=request)
return dispute
@router.post("/operator/disputes/{dispute_id}/resolve", summary="争议调解/解决")
async def resolve_dispute(
dispute_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("operator")),
):
dispute = await db.disputes.set_status(dispute_id, "resolved", resolution="平台调解结案")
if dispute is None:
raise HTTPException(status_code=404, detail="Dispute not found")
await write_audit(db, action="dispute.resolve", resource="dispute", resource_id=dispute_id,
user=actor, request=request)
return dispute
# ── 投资撮合(按偏好排序项目)───────────────────────────────────────────────
@router.get("/investor/matches", summary="投资撮合(按偏好推荐)")
async def investor_matches(
db: Database = Depends(get_db),
user: dict = Depends(require_roles("investor")),
):
pref = await db.investor_prefs.get(user["id"]) or {}
industries = set(pref.get("industries") or [])
data = await db.portal_pages.get("investor", "projects") or {"items": []}
items = data.get("items", [])
scored = []
for it in items:
text = f"{it.get('title', '')} {it.get('meta', '')}"
score = 0
for ind in industries:
if ind and ind in text:
score += 10
stage = pref.get("stage")
if stage and stage in text:
score += 5
scored.append((score, it))
scored.sort(key=lambda x: x[0], reverse=True)
return {"items": [{"title": it["title"], "meta": it["meta"], "tag": it.get("tag"),
"match": min(100, s + 50)} for s, it in scored]}
# ── 路演参与(进入直播/出席)────────────────────────────────────────────────
@router.post("/roadshows/{roadshow_id}/join", summary="进入路演直播")
async def join_roadshow(
roadshow_id: str,
request: Request,
db: Database = Depends(get_db),
actor: dict = Depends(require_roles("investor", "government", "operator", "carrier", "enterprise")),
):
rs = await db.roadshows.get(roadshow_id)
if rs is None:
raise HTTPException(status_code=404, detail="Roadshow not found")
# 标记出席(若已报名则更新为 attended;否则记录出席)
await db.roadshow_regs.create(roadshow_id, actor["id"], role="attendee")
await write_audit(db, action="roadshow.join", resource="roadshow", resource_id=roadshow_id,
user=actor, request=request)
return {"joined": True, "live_url": rs.get("live_url") or "http://live.example/roadshow"}