374 lines
15 KiB
Python
374 lines
15 KiB
Python
# -*- coding: utf-8 -*-
|
|
"""横切能力端点:通知/信用/合同/结算/争议/评价/投资撮合/路演参与。
|
|
"""
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import json
|
|
|
|
from fastapi import APIRouter, Depends, HTTPException, Request
|
|
from pydantic import BaseModel
|
|
|
|
from ..dependencies import get_db
|
|
from ..schemas.ecosystem import RatingRequest, SignRequest, DisputeRequest
|
|
from ..schemas.operator import ActivityCreateRequest, ActivityUpdateRequest
|
|
from ...rbac import require_permission, require_roles, write_audit
|
|
from ...infrastructure.repositories import Database
|
|
from ...services import training_admin_bridge as training
|
|
|
|
router = APIRouter(tags=["ecosystem"])
|
|
|
|
_ANY = ("opc_member", "carrier", "operator")
|
|
|
|
|
|
# ── 通知中心(实时通知体系)──────────────────────────────────────────────
|
|
@router.get("/notifications", summary="我的通知(分页/分类/已读筛选)")
|
|
async def list_notifications(
|
|
db: Database = Depends(get_db),
|
|
user: dict = Depends(require_roles(*_ANY)),
|
|
category: str = "",
|
|
read: str = "",
|
|
page: int = 1,
|
|
page_size: int = 50,
|
|
):
|
|
read_flag = None if read == "" else (read == "1" or read.lower() == "true")
|
|
res = await db.notifications.list_for(
|
|
user["id"], limit=max(1, min(page_size, 100)),
|
|
category=category, read=read_flag, page=page)
|
|
return {"items": res["items"], "total": res["total"],
|
|
"unread": await db.notifications.unread(user["id"]),
|
|
"unread_by_category": await db.notifications.unread_by_category(user["id"]),
|
|
"categories": await db.notifications.all_by_category(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.post("/notifications/read/{notification_id}", summary="单条已读")
|
|
async def read_one_notification(
|
|
notification_id: str,
|
|
db: Database = Depends(get_db),
|
|
user: dict = Depends(require_roles(*_ANY)),
|
|
):
|
|
ok = await db.notifications.mark_read_one(user["id"], notification_id)
|
|
return {"ok": ok}
|
|
|
|
|
|
@router.get("/notifications/unread-count", summary="未读计数(角标轮询/首屏)")
|
|
async def unread_count(
|
|
db: Database = Depends(get_db),
|
|
user: dict = Depends(require_roles(*_ANY)),
|
|
):
|
|
return {"total": await db.notifications.unread(user["id"]),
|
|
"by_category": await db.notifications.unread_by_category(user["id"])}
|
|
|
|
|
|
@router.get("/announcements", summary="平台公告(广播历史)")
|
|
async def list_announcements(
|
|
db: Database = Depends(get_db),
|
|
user: dict = Depends(require_roles(*_ANY)),
|
|
limit: int = 20,
|
|
):
|
|
return {"items": await db.announcements.list_recent(limit=max(1, min(limit, 50)))}
|
|
|
|
|
|
# ── 通知实时通道:SSE(MQTT 不可用时前端回退;token 鉴权,按 user 过滤)────
|
|
@router.get("/notify/events")
|
|
async def notify_events(request: Request, db: Database = Depends(get_db)):
|
|
from fastapi.responses import StreamingResponse
|
|
|
|
from ...jwt import decode_access_token
|
|
from ...park.event_bus import bus
|
|
|
|
token = request.query_params.get("token", "")
|
|
auth = request.headers.get("Authorization", "")
|
|
if auth.startswith("Bearer "):
|
|
token = auth[7:]
|
|
claims = decode_access_token(token) if token else None
|
|
if not claims:
|
|
from fastapi import HTTPException as _HTTPException
|
|
|
|
raise _HTTPException(status_code=401, detail="unauthorized")
|
|
uid = claims.get("sub") or claims.get("user_id") or ""
|
|
|
|
async def gen():
|
|
q = bus.subscribe()
|
|
try:
|
|
while True:
|
|
if await request.is_disconnected():
|
|
break
|
|
try:
|
|
data = await asyncio.wait_for(q.get(), 15)
|
|
except Exception: # noqa: BLE001 — 超时发 keepalive
|
|
yield ": keepalive\n\n"
|
|
continue
|
|
try:
|
|
payload = json.loads(data)
|
|
except Exception: # noqa: BLE001
|
|
yield f"data: {data}\n\n"
|
|
continue
|
|
# 只推当前用户相关(user_id == 本人 或 * 广播)
|
|
if payload.get("user_id") in (uid, "*"):
|
|
yield f"data: {json.dumps(payload, ensure_ascii=False)}\n\n"
|
|
finally:
|
|
bus.unsubscribe(q)
|
|
|
|
return StreamingResponse(gen(), media_type="text/event-stream")
|
|
|
|
|
|
# ── 信用/评价 ──────────────────────────────────────────────────────────────
|
|
@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}
|
|
|
|
|
|
@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("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,
|
|
ip=request.client.host if request.client else "")
|
|
# 信用事件:被评方按评分加分/减分(好评 +5/+10/+15,差评 -10/-20,按任务去重)
|
|
from ...services.credit_engine import on_credit_event
|
|
score = int(req.score)
|
|
event_code = "task.review.positive" if score >= 3 else "task.review.negative"
|
|
await on_credit_event(
|
|
db, user_id=req.to_id, event_code=event_code,
|
|
payload={"score": score}, ref_type="task", ref_id=task_id,
|
|
reason=f"任务互评 {score} 星",
|
|
)
|
|
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("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
|
|
|
|
|
|
@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("opc_member")),
|
|
):
|
|
from ...services.settlement_service import SettlementService
|
|
|
|
contract = await SettlementService(db).sign_contract(task_id, actor, req.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("opc_member")),
|
|
):
|
|
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['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)}
|
|
|
|
|
|
# ── 争议 ──────────────────────────────────────────────────────────────────
|
|
@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("opc_member")),
|
|
):
|
|
from ...services.settlement_service import SettlementService
|
|
|
|
dispute = await SettlementService(db).create_dispute(task_id, actor, req.reason)
|
|
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")),
|
|
):
|
|
from ...services.settlement_service import SettlementService
|
|
|
|
dispute = await SettlementService(db).resolve_dispute(dispute_id)
|
|
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("opc_member")),
|
|
):
|
|
from ...services.settlement_service import SettlementService
|
|
|
|
return await SettlementService(db).investor_matches(user["id"])
|
|
|
|
|
|
# ── 路演参与(进入直播/出席)────────────────────────────────────────────────
|
|
@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("opc_member", "operator", "carrier")),
|
|
):
|
|
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"}
|
|
|
|
|
|
# ── 活动(成员端:OPC 用户查看 / 发布 / 报名 / 管理自己发布的)────────────
|
|
|
|
class BookingCreateRequest(BaseModel):
|
|
name: str = ""
|
|
contact: str = ""
|
|
want: str = ""
|
|
|
|
|
|
@router.get("/activities", summary="活动列表(OPC 用户)")
|
|
async def list_activities_public(
|
|
status: str = "",
|
|
user: dict = Depends(require_roles(*_ANY)),
|
|
):
|
|
return training.list_activities_public(user_id=user["id"], status=status or None)
|
|
|
|
|
|
@router.get("/activities/mine/registered", summary="我报名的活动列表")
|
|
async def my_registered_activities(
|
|
user: dict = Depends(require_roles(*_ANY)),
|
|
):
|
|
return training.list_my_registered(user["id"])
|
|
|
|
|
|
@router.get("/activities/{event_id}", summary="活动详情(发布者可看报名列表)")
|
|
async def get_activity_detail(
|
|
event_id: str,
|
|
user: dict = Depends(require_roles(*_ANY)),
|
|
):
|
|
ev = training.get_event_detail(user["id"], event_id)
|
|
if ev is None:
|
|
raise HTTPException(status_code=404, detail="活动不存在")
|
|
return ev
|
|
|
|
|
|
@router.post("/activities", summary="发布活动")
|
|
async def create_activity_public(
|
|
req: ActivityCreateRequest,
|
|
request: Request,
|
|
db: Database = Depends(get_db),
|
|
actor: dict = Depends(require_roles(*_ANY)),
|
|
):
|
|
payload = req.model_dump(exclude_none=True)
|
|
payload["publisher_id"] = actor["id"]
|
|
payload["publisher_name"] = actor.get("username") or actor.get("nickname") or ""
|
|
payload["review_status"] = "approved" # 演示环境直接通过;如需审核可改为 pending
|
|
activity = training.create_activity(payload)
|
|
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}/register", summary="报名活动")
|
|
async def register_activity_public(
|
|
event_id: str,
|
|
req: BookingCreateRequest,
|
|
request: Request,
|
|
db: Database = Depends(get_db),
|
|
actor: dict = Depends(require_roles(*_ANY)),
|
|
):
|
|
booking = training.register_event(event_id, actor["id"], actor.get("username") or "", req.model_dump(exclude_none=True))
|
|
if booking is None:
|
|
raise HTTPException(status_code=404, detail="活动不存在")
|
|
await write_audit(db, action="booking.create", resource="booking", resource_id=booking["booking"]["id"],
|
|
detail=event_id, user=actor, request=request)
|
|
return booking
|
|
|
|
|
|
@router.put("/activities/{event_id}", summary="编辑自己发布的活动")
|
|
async def update_activity_owned(
|
|
event_id: str,
|
|
req: ActivityUpdateRequest,
|
|
request: Request,
|
|
db: Database = Depends(get_db),
|
|
actor: dict = Depends(require_roles(*_ANY)),
|
|
):
|
|
activity = training.update_activity_owned(
|
|
event_id, req.model_dump(exclude_unset=True, exclude_none=True), actor["id"],
|
|
)
|
|
if activity is None:
|
|
raise HTTPException(status_code=404, detail="活动不存在或无权修改")
|
|
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_owned(
|
|
event_id: str,
|
|
request: Request,
|
|
db: Database = Depends(get_db),
|
|
actor: dict = Depends(require_roles(*_ANY)),
|
|
):
|
|
ok = training.delete_activity_owned(event_id, actor["id"])
|
|
if not ok:
|
|
raise HTTPException(status_code=404, detail="活动不存在或无权删除")
|
|
await write_audit(db, action="activity.delete", resource="activity", resource_id=event_id,
|
|
detail=event_id, user=actor, request=request)
|
|
return {"ok": True} |