Files
server-core/app/park/routers.py
T

1333 lines
51 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 -*-
"""园区子应用(app/park)—— 大屏控制 / 数据 / 设置 / 园区企业(骨架版)。
迁自 park-desktop/backend/app/routers.py,本骨架仅含不依赖重模块(llm/rag/asr/
s2s/vision)的大屏控制端点;AI/知识库/语音/视觉在后续全量迁移补齐。
全部端点经 dispatcher `/park` 前缀暴露(include_router 挂 `/park`)。
"""
from __future__ import annotations
import json
import logging
import os
import time
from datetime import datetime
from pathlib import Path
from fastapi import APIRouter, Depends, File, Header, HTTPException, Request, UploadFile
from fastapi.responses import FileResponse, JSONResponse, StreamingResponse
from pydantic import BaseModel, Field
from sqlalchemy.ext.asyncio import AsyncSession
from ..infrastructure.db import get_session
from ..infrastructure.repositories import TaskRepository
from ..jwt import decode_access_token
from ..rbac import write_audit
from ..api.schemas.admin import ParkAdmissionReviewRequest, ParkTransferReviewRequest
from . import park_config, tenants
from .auth import create_device_token, create_token, parse_device_token, parse_token, require_tenant
from .config import settings
from .event_bus import bus
from .mqtt import hub
from .sim_engine import engine_data, get_engine, refresh_engine, sim_engine
from .storage import storage
logger = logging.getLogger("dpm.api")
async def _resolve_tenant(authorization: str | None, tenant_id: str | None) -> str:
"""解析目标园区:?tenant_id=admin 亮传)> 大屏 tenant-token > 设备 token > 默认园区。"""
if tenant_id:
return tenant_id
if authorization and authorization.startswith("Bearer "):
tok = authorization.split(" ", 1)[1]
try:
return parse_token(tok) # 大屏 tenant token
except HTTPException:
pass
try:
return parse_device_token(tok)["tenant_id"] # 绑定后的设备 token
except HTTPException:
pass
return await tenants.ensure_default_tenant()
log = logging.getLogger("dpm.api")
router = APIRouter()
# ==================== 园区认证 + 租户管理 ====================
class LoginBody(BaseModel):
username: str
password: str
class TenantBody(BaseModel):
name: str = "新园区"
intro: list[str] = ["", ""]
username: str = "admin"
password: str = "123456"
class TenantPatch(BaseModel):
name: str | None = None
intro: list[str] | None = None
username: str | None = None
password: str | None = None
class BindBody(BaseModel):
username: str
@router.post("/auth/login", summary="园区端登录(绑定为园区管理员的平台账号)")
async def park_login(body: LoginBody):
# 登录 = 校验平台账号密码 + 该账号已绑定为某园区管理员;成功 → 签发该园区 tenant token
from app.infrastructure.repositories import Database
db = Database()
try:
user = await db.users.get_by_username(body.username.strip())
if not user or not await db.users.verify_password(user, body.password):
return JSONResponse({"ok": False, "error": "账号或密码不正确"}, status_code=401)
t = await tenants.find_by_admin(body.username.strip())
if not t:
return JSONResponse({"ok": False, "error": "该账号未绑定任何园区管理员"}, status_code=401)
if t.get("status") == "disabled":
return JSONResponse({"ok": False, "error": "该园区已禁用,禁止登录使用"}, status_code=403)
finally:
await db.close()
return {"ok": True, "tenant_id": t["id"], "name": t["name"], "intro": t["intro"], "token": create_token(t["id"])}
@router.get("/auth/me", summary="当前园区资料(页眉)")
async def park_me(authorization: str | None = Header(default=None)):
tid = require_tenant(authorization)
t = await tenants.get_tenant(tid)
if t is None:
raise HTTPException(status_code=404, detail="园区不存在")
return {"ok": True, "tenant_id": tid, "name": t["name"], "intro": t.get("intro", [])}
# ---- 租户管理(运营端)----
@router.get("/tenants", summary="园区列表")
async def tenant_list():
return {"ok": True, "tenants": await tenants.list_tenants()}
@router.post("/tenants", summary="创建园区(名称/两段式简介/账密,seed 默认数据)")
async def tenant_create(body: TenantBody):
t = await tenants.create_tenant(body.name, body.intro, body.username, body.password)
return {"ok": True, "tenant": t}
@router.put("/tenants/{tid}", summary="更新园区(名称/简介/管理员账号/密码)")
async def tenant_update(tid: str, body: TenantPatch):
patch = {k: v for k, v in body.model_dump(exclude_none=True).items() if v is not None}
t = await tenants.update_tenant(tid, patch)
if t is None:
return JSONResponse({"ok": False, "error": "园区不存在"}, status_code=404)
refresh_engine(tid)
return {"ok": True, "tenant": t}
@router.delete("/tenants/{tid}", summary="删除园区")
async def tenant_delete(tid: str):
ok = await tenants.delete_tenant(tid)
return {"ok": ok, "id": tid}
@router.post("/tenants/{tid}/bind", summary="绑定园区管理员(已有平台账号)")
async def tenant_bind(tid: str, body: BindBody):
from app.infrastructure.repositories import Database
db = Database()
try:
user = await db.users.get_by_username(body.username.strip())
user_id = (user or {}).get("id", "")
finally:
await db.close()
if user is None:
return JSONResponse({"ok": False, "error": "账号不存在"}, status_code=404)
ok = await tenants.bind_admin(tid, body.username.strip(), operator_user_id=user_id)
# 绑定即成为该园区的载体方(carrier):账号单一角色,由后台绑定设置(登录后按 carrier 进园区工作台)。
await db.users.set_role(user_id, "carrier", None, None, None)
return {"ok": ok, "tenant_id": tid, "admin_username": body.username.strip()}
@router.post("/tenants/{tid}/unbind", summary="解除园区管理员")
async def tenant_unbind(tid: str):
ok = await tenants.unbind_admin(tid)
return {"ok": ok, "tenant_id": tid}
class StatusBody(BaseModel):
status: str # active | disabled
@router.post("/tenants/{tid}/status", summary="禁用/启用园区(不删,禁用禁止登录)")
async def tenant_status(tid: str, body: StatusBody):
ok = await tenants.set_status(tid, body.status)
if not ok:
return JSONResponse({"ok": False, "error": "无效状态或园区不存在"}, status_code=400)
return {"ok": True, "tenant_id": tid, "status": body.status}
# ==================== 数据模型 ====================
class SettingsBody(BaseModel):
volume: int | None = None
sfx_volume: int | None = None
play_mode: str | None = None
image_duration: int | None = None
fullscreen: bool | None = None
autostart: bool | None = None
username: str | None = None
password: str | None = None
class ActionBody(BaseModel):
action: str
class PathBody(BaseModel):
path: str
class UrlBody(BaseModel):
url: str
name: str = ""
type: str = "image"
class DisplayCommandBody(BaseModel):
action: str
params: dict = Field(default_factory=dict)
screen_id: str = ""
screen_role: str = ""
class RegisterBody(BaseModel):
device_id: str = ""
role: str = ""
parent_id: str = ""
class CompanyBody(BaseModel):
name: str = ""
zone: str = ""
room: str = ""
industry: str = ""
bio: str = ""
founder: str = ""
status: str = "applying"
employees: int | None = None
owner_user_id: str = ""
legal_person: str = ""
legal_phone: str = ""
registered_capital: str = ""
company_type: str = ""
honors: str = ""
address: str = ""
contact_phone: str = ""
founded_at: str = ""
emp_total: int | None = None
emp_grad: int | None = None
emp_layoff: int | None = None
emp_veteran: int | None = None
emp_migrant: int | None = None
class AgentBody(BaseModel):
system_prompt: str | None = None
model: str | None = None
enabled_tools: list[str] | None = None
preset_questions: list[str] | None = None
opening: str | None = None
class KbDocBody(BaseModel):
group: str = "general"
title: str = ""
content_md: str = ""
class ScreenBody(BaseModel):
name: str | None = None
region: str | None = None
capacity: int | None = None
invested: int | None = None
jobs: int | None = None
area: int | None = None
founded: int | None = None
address: str | None = None
phone: str | None = None
email: str | None = None
revenue_total: float | None = None
revenue_tax: float | None = None
intro: list[str] | None = None
feed: list[dict] | None = None
zones: list[dict] | None = None
industryMix: list[dict] | None = None
# ==================== 健康 / 设置 / 引导 ====================
@router.get("/api/health")
async def health():
return {
"ok": True,
"service": "opc-park",
"mqtt_connected": hub.connected,
"screens_online": hub.screens_online(),
"version": "1.0.0",
}
# ==================== 媒体上传 / 播放控制 / 客户端更新(迁自 park-desktop backend ====================
@router.post("/api/upload")
async def upload_media(file: UploadFile = File(...), tenant_id: str | None = None):
"""上传媒体(图片/视频)到园区媒体目录,返回 {path, url}。"""
from .config import settings
fname = Path(file.filename or "upload.bin").name
dest = settings.MEDIA_DIR / fname
dest.parent.mkdir(parents=True, exist_ok=True)
data = await file.read()
dest.write_bytes(data)
return {"ok": True, "path": f"/park/file/{fname}", "url": f"/park/file/{fname}"}
@router.get("/api/state")
async def get_state():
return {"status": "idle", "index": 0, "name": "", "media_type": ""}
@router.post("/api/control")
async def control(body: ActionBody):
cmd = hub.publish_command(body.action)
return {"ok": cmd.get("published", False)}
_INSTALLER_DIR = Path(__file__).resolve().parent.parent.parent / "serverdata" / "park" / "installers"
@router.get("/download")
async def download_installer():
"""下载 Windows 客户端安装包(serverdata/park/installers/ 下最新 *.setup.exe)。"""
files = sorted(_INSTALLER_DIR.glob("*setup.exe"), reverse=True)
if not files:
return JSONResponse({"ok": False, "error": "安装包不存在"}, status_code=404)
f = files[0]
return FileResponse(path=str(f), filename=f.name, media_type="application/octet-stream")
class UpdateCheckBody(BaseModel):
current: str = ""
@router.post("/api/update/check")
async def update_check(body: UpdateCheckBody):
"""客户端更新检测:返回最新版本号、安装包下载地址、是否需要更新。"""
files = sorted(_INSTALLER_DIR.glob("*setup.exe"), reverse=True)
if not files:
return {"ok": False, "error": "installer not found"}
install = files[0].name
import re
m = re.search(r"_(\d+\.\d+(?:\.\d+)?)", install)
latest = m.group(1) if m else ""
current = (body.current or "").strip()
update_available = bool(latest) and (not current or current != latest)
return {"ok": True, "latest_version": latest, "installer_name": install, "download_url": "/download", "update_available": update_available}
# ==================== 媒体 / 播放列表(storage → serverdata/park/data.json ====================
def _media_files() -> list[dict]:
from .config import settings
out = []
for f in sorted(settings.MEDIA_DIR.glob("*")):
if f.is_file() and f.suffix.lower() in {".mp4", ".mkv", ".avi", ".jpg", ".jpeg", ".png"}:
out.append({"relative_path": f.name, "name": f.name,
"type": "video" if f.suffix.lower() in {".mp4", ".mkv", ".avi"} else "image",
"size": f.stat().st_size})
return out
@router.get("/media")
async def list_media():
return _media_files()
@router.get("/api/playlist")
async def get_playlist():
from .storage import storage
return {"playlist": storage.get_playlist()}
@router.post("/api/playlist/add")
async def add_playlist(body: PathBody):
from .storage import storage
storage.add_to_playlist(body.path)
hub.publish_command("playlist_changed", {"action": "add", "path": body.path})
return {"ok": True}
@router.post("/api/playlist/remove")
async def remove_playlist(body: PathBody):
from .storage import storage
storage.remove_from_playlist(body.path)
hub.publish_command("playlist_changed", {"action": "remove", "path": body.path})
return {"ok": True}
@router.post("/api/playlist/play")
async def play_playlist(body: PathBody):
hub.publish_command("play_target", {"path": body.path})
return {"ok": True}
@router.post("/api/media/add-url")
async def add_url(body: UrlBody):
from .storage import storage
created = storage.add_url_media(body.url, body.name, body.type)
if not created:
return {"ok": False, "error": "已存在"}
return {"ok": True}
@router.post("/api/delete")
async def delete_media(body: PathBody):
from .storage import storage
storage.delete_media(body.path)
return {"ok": True}
@router.get("/api/settings")
async def get_settings():
return storage.get_settings()
@router.post("/api/settings")
async def update_settings(body: SettingsBody):
patch = body.model_dump(exclude_none=True)
storage.update_settings(**patch)
hub.publish_command("settings_changed", patch)
return storage.get_settings()
@router.get("/api/config")
async def runtime_config(request: Request):
"""大屏启动引导:下发 MQTT/API/语音地址(顶层 mqtt_url/mqtt_username/mqtt_password 字段,兼容大屏 bootstrap)。"""
return {
"mqtt": {
"host": settings.MQTT_HOST,
"port": settings.MQTT_PORT,
"ws_url": settings.MQTT_WS_URL,
"tick_interval": settings.MQTT_TICK_INTERVAL,
},
# 大屏 main.jsx bootstrap 读顶层字段;缺省时回退构建期烘焙地址
"mqtt_url": settings.MQTT_WS_URL,
"mqtt_username": settings.MQTT_USERNAME or "",
"mqtt_password": settings.MQTT_PASSWORD or "",
"voice_ws": settings.S2S_WS_URL if settings.S2S_ENABLED else "",
"media_base": f"{request.base_url}park/file",
# 大屏前端 request() 以 `${api_base}/api/...`、`${api_base}/file/...` 拼接,故 api_base 返回 `/park` 前缀
"api_base": f"{request.base_url}park",
}
class VisionConfigBody(BaseModel):
api: str = ""
@router.get("/api/vision/config")
async def vision_config_get(authorization: str | None = Header(None), tenant_id: str | None = None):
"""园区端读:本园区视频识别服务地址(空=禁用)。识别能力在独立 park-vision。"""
tid = await _resolve_tenant(authorization, tenant_id)
api = await tenants.get_tenant_vision(tid)
return {"ok": True, "api": api}
@router.put("/api/vision/config")
async def vision_config_set(body: VisionConfigBody, authorization: str | None = Header(None), tenant_id: str | None = None):
"""园区端写:配置本园区视频识别服务地址(可填局域网地址;空=禁用)。"""
tid = await _resolve_tenant(authorization, tenant_id)
api = (body.api or "").strip().rstrip("/")
await tenants.set_tenant_vision(tid, api)
return {"ok": True, "api": api}
# ==================== 大屏数据(sim_engine ====================
@router.get("/api/dashboard/snapshot")
async def dashboard_snapshot(authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
return get_engine(tid).snapshot()
@router.get("/api/dashboard/overview")
async def dashboard_overview(authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
return get_engine(tid).snapshot()
@router.get("/api/tasks", summary="大屏任务展示(平台全局已发布任务)")
async def park_tasks(
authorization: str | None = Header(None),
tenant_id: str | None = None,
session: AsyncSession = Depends(get_session),
):
"""返回系统任务中心的上墙任务(published/claimed/doing),供大屏轮播卡片。
大屏卡片含任务 ID、task_code 与二维码(scan_payload)。任务为平台全局实体,
不受园区隔离;此处按园区身份校验后返回。
"""
tid = await _resolve_tenant(authorization, tenant_id)
repo = TaskRepository(session)
# 大屏仅展示「当前园区可接」:本园专属(park_id=本园) 公有(park_public),状态已发布/在接/在做。
items = [t for t in await repo.list_published()
if (t.get("park_id") == tid) or t.get("park_public")]
now_label = datetime.now().strftime("%Y%m%d")
def payload(t: dict) -> dict:
return {
"id": t["id"],
"task_code": t["task_code"] or f"TK-{now_label}-{t['id'][-5:]}",
"scan_payload": t["task_code"] or f"TK-{now_label}-{t['id'][-5:]}",
"title": t["title"],
"category": t["category"],
"summary": (t["description"] or "")[:160],
"budget_min": t["budget_min"],
"budget_max": t["budget_max"],
"mode": t.get("mode", "grab"),
"deadline": t.get("deadline", ""),
"delivery_days": t.get("delivery_days", 0),
"headcount": t.get("headcount", 0),
"exclusive": t.get("exclusive", False),
"publisher_name": t.get("publisher_name", ""),
"status": t["status"],
"claimed_by": t["claimed_by"],
"claimed_at": t["claimed_at"],
"tags": [tag for tag in (t.get("tags") or "").split(",") if tag],
}
return {"items": [payload(t) for t in items]}
class ParkReleaseBody(BaseModel):
vis: str = "public" # 发单后可见性 public/c_visible
class ParkAssignBody(BaseModel):
user_id: str
class ParkStatusBody(BaseModel):
status: str # published/cancelled
@router.get("/api/park/tasks", summary="园区端任务(本园发布 + 指派本园)")
async def park_task_list(
authorization: str | None = Header(None),
tenant_id: str | None = None,
session: AsyncSession = Depends(get_session),
):
tid = await _resolve_tenant(authorization, tenant_id)
repo = TaskRepository(session)
items = await repo.list()
mine = [t for t in items
if (t.get("park_id") == tid)
or (t.get("assign_type") == "park" and t.get("assign_park_id") == tid)]
return {"items": mine}
@router.post("/api/park/tasks/{task_id}/status", summary="园区任务上架/下架")
async def park_task_status(
task_id: str,
body: ParkStatusBody,
authorization: str | None = Header(None),
tenant_id: str | None = None,
session: AsyncSession = Depends(get_session),
):
tid = await _resolve_tenant(authorization, tenant_id)
repo = TaskRepository(session)
task = await repo.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="任务不存在")
if not (task.get("park_id") == tid or task.get("assign_park_id") == tid):
raise HTTPException(status_code=403, detail="非本园区任务")
if body.status not in ("published", "cancelled"):
raise HTTPException(status_code=400, detail="仅支持 published/cancelled")
return await repo.set_status(task_id, body.status)
@router.post("/api/park/tasks/{task_id}/release", summary="园区内发单(指派本园区任务 → 园区企业可抢)")
async def park_task_release(
task_id: str,
body: ParkReleaseBody,
authorization: str | None = Header(None),
tenant_id: str | None = None,
session: AsyncSession = Depends(get_session),
):
tid = await _resolve_tenant(authorization, tenant_id)
repo = TaskRepository(session)
task = await repo.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="任务不存在")
if not (task.get("park_id") == tid or task.get("assign_park_id") == tid):
raise HTTPException(status_code=403, detail="非本园区任务")
if task.get("status") not in ("published", "claimed"):
raise HTTPException(status_code=400, detail="任务当前不可发单")
updated = await repo.update(task_id, {"park_released": True, "visibility": body.vis or "public"})
return updated
@router.post("/api/park/tasks/{task_id}/assign", summary="园区内指派(分派给本园某 OPC/企业)")
async def park_task_assign(
task_id: str,
body: ParkAssignBody,
authorization: str | None = Header(None),
tenant_id: str | None = None,
session: AsyncSession = Depends(get_session),
):
tid = await _resolve_tenant(authorization, tenant_id)
repo = TaskRepository(session)
task = await repo.get(task_id)
if task is None:
raise HTTPException(status_code=404, detail="任务不存在")
if not (task.get("park_id") == tid or task.get("assign_park_id") == tid):
raise HTTPException(status_code=403, detail="非本园区任务")
if not body.user_id:
raise HTTPException(status_code=400, detail="缺少被指派对象")
await repo.update(task_id, {"assign_opc_id": body.user_id, "park_released": True})
if task.get("status") == "published":
await repo.claim(task_id, body.user_id)
return await repo.get(task_id)
@router.get("/api/park/zones")
async def park_zones(authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
return get_engine(tid).snapshot().get("zones", [])
# ==================== 屏幕管理(园区端创建/管理大屏设备) ====================
class ScreenBody2(BaseModel):
device_id: str = ""
name: str = ""
role: str = "main"
location: str = ""
@router.get("/api/screens")
async def screens_list(authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
return await tenants.list_screens(tid)
@router.post("/api/screens")
async def screens_create(body: ScreenBody2, authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
return await tenants.create_screen(tid, body.model_dump())
@router.delete("/api/screens/{sid}")
async def screens_delete(sid: str, authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
ok = await tenants.delete_screen(tid, sid)
return {"ok": ok, "id": sid}
# ==================== 大屏设备注册 / 绑定(连接码 + MQTT 通知) ====================
class DeviceBody(BaseModel):
device_id: str = ""
class BindByCodeBody(BaseModel):
code: str
name: str = ""
role: str = "main"
location: str = ""
def _gen_code() -> str:
import random
return str(random.randint(10000000, 99999999))
def _publish_bind(device_id: str, payload: dict) -> None:
"""未绑定大屏仅订阅 bind 频道;绑定成功后经此频道通知并下发 token。
retain=True:绑定/解绑状态持久,屏幕即使延迟连接,订阅时也会立即收到最新状态(避免
绑定成功但屏幕 MQTT 尚未订阅导致丢失 bound 事件、无法自动进入)。"""
hub.publish(f"opc/display/bind/{device_id}", payload, qos=1, retain=True)
@router.post("/api/devices/register")
async def device_register(body: DeviceBody):
device_id = (body.device_id or "").strip()
if not device_id:
return JSONResponse({"ok": False, "error": "device_id 必填"}, status_code=400)
d = await tenants.ensure_device(device_id)
return {"ok": True, "device_id": device_id, "bound": d.get("status") == "bound", "tenant_id": d.get("tenant_id")}
@router.post("/api/devices/{device_id}/code")
async def device_code(device_id: str):
d = await tenants.ensure_device(device_id)
code = _gen_code()
expires = time.strftime("%Y-%m-%dT%H:%M:%S", time.localtime(time.time() + 600))
await tenants.set_device_code(device_id, code, expires)
return {"ok": True, "code": code, "expires": expires}
@router.post("/tenants/{tid}/screens/bind-by-code")
async def bind_screen_by_code(tid: str, body: BindByCodeBody):
"""园区端输入/扫码绑定:校验码 → 绑定到该园区 → 签发设备 token + MQTT 通知大屏。"""
d = await tenants.find_device_by_code(body.code.strip())
if d is None:
return JSONResponse({"ok": False, "error": "连接码无效或已过期"}, status_code=404)
dev_id = d["device_id"]
if d.get("tenant_id") and d["tenant_id"] != tid:
return JSONResponse({"ok": False, "error": "该大屏已绑定其它园区"}, status_code=409)
screen = await tenants.bind_device(tid, dev_id, body.name or "", body.role or "main", body.location or "")
token = create_device_token(tid, dev_id)
tenant = await tenants.get_tenant(tid)
vision_api = await tenants.get_tenant_vision(tid) # 本园区视频识别地址(可局域网),随绑定下发大屏
_publish_bind(dev_id, {"event": "bound", "token": token, "tenant_id": tid,
"name": (tenant or {}).get("name", ""), "intro": (tenant or {}).get("intro", []),
"vision_api": vision_api})
return {"ok": True, "screen": screen, "token": token, "vision_api": vision_api}
@router.post("/api/devices/{device_id}/unbind")
async def device_unbind(device_id: str, authorization: str | None = Header(None), tenant_id: str | None = None):
await _resolve_tenant(authorization, tenant_id)
ok = await tenants.unbind_device(device_id)
if ok:
# 通知大屏退出到待绑定页(大屏订阅 bind/<device> 频道,收到 unbound 即清理本地并回绑定页)
_publish_bind(device_id, {"event": "unbound"})
return {"ok": ok, "device_id": device_id}
# ==================== 入驻企业(租户 data.companies 主数据源) ====================
@router.get("/api/park/companies")
async def park_companies(authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
return await tenants.list_companies(tid)
@router.post("/api/park/companies")
async def park_company_create(body: CompanyBody, authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
if not body.owner_user_id:
return JSONResponse({"ok": False, "error": "请绑定负责人"}, status_code=400)
row = await tenants.create_company(tid, body.model_dump())
# 负责人为企业首个成员
await tenants.add_company_member(row["id"], body.owner_user_id)
return row
@router.put("/api/park/companies/{cid}")
async def park_company_update(cid: str, body: CompanyBody, authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
row = await tenants.update_company(tid, cid, body.model_dump(exclude_unset=True))
if row is None:
return JSONResponse({"ok": False, "error": "企业不存在"}, status_code=404)
return row
@router.delete("/api/park/companies/{cid}")
async def park_company_delete(cid: str, authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
ok = await tenants.delete_company(tid, cid)
return {"ok": ok, "id": cid}
# ==================== 园区智能体 / 知识库 / 大屏数据 ====================
@router.get("/api/agent/config")
async def agent_config_get(authorization: str | None = Header(None), tenant_id: str | None = None):
return await tenants.get_agent(await _resolve_tenant(authorization, tenant_id))
@router.put("/api/agent/config")
async def agent_config_put(body: AgentBody, authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
return await tenants.update_agent(tid, body.model_dump(exclude_none=True))
@router.get("/api/kb/docs")
async def kb_docs_list(authorization: str | None = Header(None), tenant_id: str | None = None):
return await tenants.kb_docs(await _resolve_tenant(authorization, tenant_id))
@router.post("/api/kb/docs")
async def kb_docs_create(body: KbDocBody, authorization: str | None = Header(None), tenant_id: str | None = None):
return await tenants.create_kb_doc(await _resolve_tenant(authorization, tenant_id), body.model_dump())
@router.put("/api/kb/docs/{did}")
async def kb_docs_update(did: str, body: KbDocBody, authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
row = await tenants.update_kb_doc(tid, did, body.model_dump(exclude_unset=True))
if row is None:
return JSONResponse({"ok": False, "error": "文档不存在"}, status_code=404)
return row
@router.delete("/api/kb/docs/{did}")
async def kb_docs_delete(did: str, authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
ok = await tenants.delete_kb_doc(tid, did)
return {"ok": ok, "id": did}
@router.post("/api/kb/reindex")
async def kb_reindex(authorization: str | None = Header(None), tenant_id: str | None = None):
# 知识库索引重建(重模块 rag 后续接入);当前仅返回 ok 占位
return {"ok": True, "reindexed": False}
@router.get("/api/screen/data")
async def screen_data_get(authorization: str | None = Header(None), tenant_id: str | None = None):
return await tenants.screen_view(await _resolve_tenant(authorization, tenant_id))
@router.put("/api/screen/data")
async def screen_data_put(body: ScreenBody, authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
return await tenants.update_screen(tid, body.model_dump(exclude_none=True))
# ==================== AI / 知识库 / 语音识别 / 视觉(重模块已迁入) ====================
class ChatMessage(BaseModel):
role: str
content: str
class ChatBody(BaseModel):
messages: list[ChatMessage]
screen_id: str = ""
class ToolsExecBody(BaseModel):
name: str = ""
args: dict = {}
screen_id: str = ""
class VisionEventBody(BaseModel):
event: str
faces: int = 0
dwell_ms: int = 0
detail: str = ""
AI_GROUPED_QUESTIONS = [
{"title": "入驻与流程", "questions": ["如何申请入驻园区?", "园区入驻的条件有哪些?", "入驻需要准备哪些材料?", "入驻流程是怎样的?", "入驻评审如何打分?", "入驻需要多长时间?"]},
{"title": "政策与扶持", "questions": ["园区有哪些创业政策扶持?", "如何申请创业补贴?", "如何申请创业担保贷款?", "对高校毕业生有什么优惠?", "科技成果转化有哪些支持?"]},
{"title": "OPC 概念", "questions": ["什么是 OPC", "OPC 创业有哪些模式?", "OPC 适合哪些人?", "OPC 创业者从哪里开始?", "OPC 常用的人工智能工具有哪些?"]},
{"title": "场地与服务", "questions": ["园区提供哪些免费办公空间?", "园区有哪些孵化服务?", "园区有哪些创业辅导?", "园区有哪些基础配套?", "园区可以免费使用哪些资源?"]},
{"title": "园区企业介绍", "questions": ["介绍一下园区入驻企业", "园区有哪些 AI 科技企业?", "园区有哪些跨境电商企业?", "园区有哪些生物医药企业?", "介绍一下云南派音人工智能科技"]},
]
AI_OPC_TOOLS = ["DeepSeek", "通义千问", "ChatGPT", "豆包", "Midjourney", "Stable Diffusion", "剪映", "Notion AI", "WPS AI", "GitHub Copilot"]
@router.get("/api/ai/questions")
async def ai_questions():
return {"groups": AI_GROUPED_QUESTIONS, "opcTools": AI_OPC_TOOLS}
@router.post("/api/ai/chat")
async def ai_chat(body: ChatBody):
from .llm import run_chat as llm_run_chat
messages = [m.model_dump() for m in body.messages]
return llm_run_chat(messages, body.screen_id or "")
@router.post("/api/ai/asr")
async def ai_asr(file: UploadFile = File(...), format: str = "m4a"):
"""语音识别:上传录音(表单 file)→ 阿里云 paraformer 转写为文本。"""
from .asr import transcribe as asr_transcribe
data = await file.read()
if not data:
return JSONResponse({"ok": False, "error": "空音频"}, status_code=400)
try:
return {"ok": True, "text": asr_transcribe(data, fmt=format)}
except Exception as e: # noqa: BLE001
return JSONResponse({"ok": False, "error": str(e)}, status_code=502)
@router.post("/api/tools/exec")
def tools_exec(body: ToolsExecBody):
from .tools import exec_tool
name = (body.name or "").strip()
if not name:
return {"ok": False, "error": "missing tool name"}
result = exec_tool(name, body.args or {}, body.screen_id or "")
return {"ok": True, "result": result}
@router.get("/api/kb/brief")
def kb_brief(max_chars: int = 1200):
try:
from .rag import brief
return {"ok": True, "brief": brief(max_chars=max_chars)}
except Exception as e: # noqa: BLE001
return {"ok": False, "error": str(e)}
@router.get("/api/kb/retrieve")
def kb_retrieve(q: str = "", top_k: int = 3):
if not q.strip():
return {"ok": False, "error": "missing q"}
try:
from .rag import retrieve
return {"ok": True, "chunks": retrieve(q, top_k=top_k)}
except Exception as e: # noqa: BLE001
return {"ok": False, "error": str(e)}
@router.get("/api/s2s/instructions")
def s2s_instructions():
try:
from .rag import build_instructions
return {"ok": True, "instructions": build_instructions()}
except Exception as e: # noqa: BLE001
return {"ok": False, "error": str(e)}
@router.post("/api/vision/event")
async def vision_event(body: VisionEventBody):
log.info("vision event=%s faces=%d dwell_ms=%d detail=%s", body.event, body.faces, body.dwell_ms, body.detail)
return {"ok": True}
# ==================== 展示控制(管理端 → MQTT ====================
# 全部 MQTT 前端控制命令白名单(与 park-desktop docs/mqtt-commands.md 一致)
_VALID_DISPLAY_ACTIONS = {
"navigate", "navigate_rel",
"vision_set",
"ai_input", "ai_preset", "ai_company", "ai_zone",
"voice_start", "voice_stop", "voice_refresh",
"play", "pause", "next", "prev", "set_mode", "play_target",
"dual_screen",
"alert", "show_card", "minimize", "settings_changed", "playlist_changed",
}
@router.post("/api/display/register")
async def display_register(body: RegisterBody, authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
hub.register_device(body.device_id, body.role, body.parent_id, tid)
return {"ok": True, "device_id": body.device_id, "role": body.role, "parent_id": body.parent_id, "tenant_id": tid}
@router.post("/api/display/command")
async def display_command_publish(body: DisplayCommandBody, authorization: str | None = Header(None), tenant_id: str | None = None):
tid = tenant_id or require_tenant(authorization)
if body.action not in _VALID_DISPLAY_ACTIONS:
return JSONResponse({"ok": False, "error": f"不支持的指令: {body.action}"}, status_code=400)
cmd = hub.publish_command(body.action, body.params, body.screen_id or None, body.screen_role or None, tid)
return {"ok": cmd.get("published", False), "cmd_id": cmd["cmd_id"], "action": body.action,
"screen_id": body.screen_id or "", "screen_role": body.screen_role or "both",
"tenant_id": tid, "mqtt_connected": hub.connected}
@router.get("/api/display/state")
async def display_state(authorization: str | None = Header(None), tenant_id: str | None = None):
tid = await _resolve_tenant(authorization, tenant_id)
st = hub.status()
# 只返回该租户在线设备
st["devices"] = [d for d in st.get("devices", []) if _device_tenant(d) == tid]
st["screens_online"] = len(st["devices"])
st["tenant_id"] = tid
return st
def _device_tenant(device: dict) -> str:
# status() 未带 tenant_id,从 hub 设备表补齐
for cid, v in hub.devices.items():
if cid == device.get("device_id"):
return v.get("tenant_id", "")
return ""
# ==================== SSE 兼容通道(MQTT 不可用时前端回退) ====================
@router.get("/api/events")
async def sse_events(request: Request):
async def gen():
q = bus.subscribe()
try:
while True:
if await request.is_disconnected():
break
try:
data = await asyncio_timeout(q.get(), 15)
yield f"data: {data}\n\n"
except Exception: # noqa: BLE001
yield ": keepalive\n\n"
finally:
bus.unsubscribe(q)
return StreamingResponse(gen(), media_type="text/event-stream")
def asyncio_timeout(awaitable, seconds):
"""极简超时包装:避免依赖 asyncio.timeout3.11+ 的上下文管理器)。"""
import asyncio
return asyncio.wait_for(awaitable, timeout=seconds)
# ==================== 园区端(carrier)入园 / 转园 审核 —— 仅本园区 ====================
async def _carrier_db():
"""园区端审核端点统一走平台 Database(与 park_login 一致)。"""
from app.infrastructure.repositories import Database
db = Database()
try:
yield db
finally:
await db.close()
async def _carrier_user(request: Request) -> dict:
"""园区端鉴权(平台一套):解析平台 Bearer 登录令牌,要求业务角色 carrier(载体方)。
park 子应用挂在 dispatcher 下,request.app.state.db 未设置,故不能用平台 require_roles
(它读 app.state.db);这里直接 decode_access_token 验证 + role=carrier。
园区端管理员以平台 carrier 身份登录进入,不另立 park-token 鉴权。
"""
auth = request.headers.get("authorization", "")
if not auth.startswith("Bearer "):
raise HTTPException(status_code=401, detail="缺少登录令牌")
payload = decode_access_token(auth.split(" ", 1)[1])
if not payload:
raise HTTPException(status_code=401, detail="令牌无效或已过期")
if payload.get("role") != "carrier":
raise HTTPException(status_code=403, detail="仅园区载体方可访问")
return {"id": payload.get("sub", ""), "username": payload.get("username", ""), "role": "carrier"}
async def _my_park(user: dict) -> dict:
"""当前 carrier 账号关联的园区,未关联则拒。
主键 operator_user_id(绑定园区管理员时写入);兜底用 username 匹配
park_tenants.admin_username(历史/种子创建的园区管理员,operator_user_id 未回填)。
"""
t = await tenants.find_by_operator_user_id(user["id"])
if not t and user.get("username"):
t = await tenants.find_by_admin(user["username"])
if not t:
raise HTTPException(status_code=403, detail="该账号未关联任何园区")
return t
@router.get("/api/my-tenant", summary="当前 carrier 绑定园区")
async def my_tenant(
db=Depends(_carrier_db),
user: dict = Depends(_carrier_user),
):
t = await _my_park(user)
return {"ok": True, "tenant": {"id": t["id"], "name": t.get("name", ""), "intro": t.get("intro", [])}}
@router.get("/api/park-admissions", summary="园区入驻申请(本园区, 平台 carrier 身份)")
async def carrier_park_admissions(
status: str = "",
db=Depends(_carrier_db),
user: dict = Depends(_carrier_user),
):
t = await _my_park(user)
return {"items": await db.park_admissions.list(status or None, tenant_id=t["id"])}
@router.post("/api/park-admissions/{aid}/review", summary="园区入驻申请审核(本园区)")
async def carrier_park_admission_review(
aid: str,
body: ParkAdmissionReviewRequest,
request: Request,
db=Depends(_carrier_db),
user: dict = Depends(_carrier_user),
):
t = await _my_park(user)
adm = await db.park_admissions.get(aid)
if adm is None:
raise HTTPException(status_code=404, detail="入驻申请不存在")
if adm.get("tenant_id") != t["id"]:
raise HTTPException(status_code=403, detail="仅可审核本园区申请")
st = body.status
if st not in ("approved", "rejected", "reviewing"):
raise HTTPException(status_code=400, detail="状态仅支持 approved/rejected/reviewing")
updated = await db.park_admissions.set_status(aid, st,
reviewer=user.get("id", ""), comment=body.comment)
if st == "approved" and adm.get("user_id"):
await db.users.set_classification(adm["user_id"], affiliation="park",
park_id=adm.get("tenant_id", ""), park_name=adm.get("tenant_name", ""))
await write_audit(db, action="park.admission_review", resource="park_admission",
resource_id=aid, detail=f"status={st}", user=user, request=request)
return {"ok": True, "admission": updated}
@router.get("/api/park-transfers", summary="OPC 转园申请(本园区相关, 平台 carrier 身份)")
async def carrier_park_transfers(
status: str = "",
db=Depends(_carrier_db),
user: dict = Depends(_carrier_user),
):
t = await _my_park(user)
return {"items": await db.park_transfers.list(status or None, park_id=t["id"])}
@router.post("/api/park-transfers/{tid}/review", summary="OPC 转园申请审核(本园区)")
async def carrier_park_transfer_review(
tid: str,
body: ParkTransferReviewRequest,
request: Request,
db=Depends(_carrier_db),
user: dict = Depends(_carrier_user),
):
t = await _my_park(user)
tr = await db.park_transfers.get(tid)
if tr is None:
raise HTTPException(status_code=404, detail="转园申请不存在")
if tr.get("from_park_id") != t["id"] and tr.get("to_park_id") != t["id"]:
raise HTTPException(status_code=403, detail="仅可审核与本园区相关的转园")
st = body.status
if st not in ("approved", "rejected", "reviewing"):
raise HTTPException(status_code=400, detail="状态仅支持 approved/rejected/reviewing")
updated = await db.park_transfers.set_status(tid, st,
reviewer=user.get("id", ""), comment=body.comment)
if st == "approved":
await db.users.set_classification(tr["user_id"], affiliation="park",
park_id=tr.get("to_park_id", ""), park_name=tr.get("to_park_name", ""))
await write_audit(db, action="park.transfer_review", resource="park_transfer",
resource_id=tid, detail=f"status={st}", user=user, request=request)
return {"ok": True, "transfer": updated}
# ==================== 园区企业:成员 + 算力(carrier 本园区) ====================
class CompanyMemberBody(BaseModel):
company_id: str
user_id: str
class CompanyComputeBody(BaseModel):
discount: int | None = Field(default=None, ge=0, le=100)
quota_add: int = Field(default=0, ge=0)
async def _own_company(t: dict, company_id: str) -> str:
"""校验企业归属本园区,返回 company 数据;否则 403。"""
c = await tenants.get_company(company_id)
if c is None:
raise HTTPException(status_code=404, detail="企业不存在")
if c.get("tenant_id") != t["id"]:
raise HTTPException(status_code=403, detail="仅可管理本园区企业")
return c
@router.get("/api/company-members", summary="园区企业成员列表")
async def carrier_company_members(company_id: str, user: dict = Depends(_carrier_user)):
t = await _my_park(user)
await _own_company(t, company_id)
return {"items": await tenants.company_members(company_id)}
@router.post("/api/company-members", summary="把用户加入企业(成员)")
async def carrier_company_add_member(body: CompanyMemberBody, user: dict = Depends(_carrier_user)):
t = await _my_park(user)
await _own_company(t, body.company_id)
r = await tenants.add_company_member(body.company_id, body.user_id)
if r is None:
raise HTTPException(status_code=404, detail="企业或用户不存在")
return r
@router.delete("/api/company-members/{company_id}/{user_id}", summary="把用户移出企业")
async def carrier_company_remove_member(company_id: str, user_id: str, user: dict = Depends(_carrier_user)):
t = await _my_park(user)
await _own_company(t, company_id)
r = await tenants.remove_company_member(company_id, user_id)
if r is None:
raise HTTPException(status_code=404, detail="企业或用户不存在")
if not r.get("ok"):
raise HTTPException(status_code=400, detail=r.get("reason", "该用户不隶属此企业"))
return r
@router.get("/api/company-pool", summary="本园区可加入企业的成员候选")
async def carrier_company_pool(user: dict = Depends(_carrier_user)):
t = await _my_park(user)
return {"items": await tenants.available_members(t["id"])}
@router.put("/api/companies/{cid}/compute", summary="设置企业算力(折扣%/配额,只增不减)+ 同步引擎")
async def carrier_company_compute(cid: str, body: CompanyComputeBody, user: dict = Depends(_carrier_user)):
from ..services import compute_client
t = await _my_park(user)
await _own_company(t, cid)
try:
updated = await tenants.set_company_compute(cid, discount=body.discount, quota_add=body.quota_add)
except ValueError as exc:
raise HTTPException(status_code=400, detail=str(exc))
except LookupError as exc:
raise HTTPException(status_code=404, detail=str(exc))
# best-effort 引擎同步:折扣 → 用户组倍率;配额 → 逐成员充值(失败不阻断本地记录)
await tenants.sync_company_engine(cid)
if body.quota_add > 0:
for m in await tenants.company_members(cid):
try:
await compute_client.grant_user_quota_by_username(m["username"], body.quota_add)
except Exception as exc: # noqa: BLE001
logging.getLogger(__name__).warning("企业配额发放失败 %s: %s", m.get("username"), exc)
return {"ok": True, "company": updated, "quota_add": body.quota_add}
class CompanyOwnerBody(BaseModel):
user_id: str
@router.put("/api/companies/{cid}/owner", summary="设置企业负责人(须为企业成员)")
async def carrier_company_set_owner(cid: str, body: CompanyOwnerBody, user: dict = Depends(_carrier_user)):
t = await _my_park(user)
await _own_company(t, cid)
r = await tenants.set_company_owner(cid, body.user_id)
if r is None:
raise HTTPException(status_code=404, detail="企业或用户不存在")
if not r.get("ok"):
raise HTTPException(status_code=400, detail=r.get("reason", "该用户不隶属此企业"))
return r
# ==================== 园区端:用户管理(本园区,新建/编辑/列) ====================
class ParkUserCreateBody(BaseModel):
username: str
password: str = ""
nickname: str = ""
gender: str = ""
birthday: str = ""
phone: str = ""
email: str = ""
company: str = ""
id_card: str = ""
ethnicity: str = ""
grad_school_major: str = ""
grad_time: str = ""
class ParkUserUpdateBody(BaseModel):
nickname: str | None = None
gender: str | None = None
birthday: str | None = None
phone: str | None = None
email: str | None = None
company: str | None = None
id_card: str | None = None
ethnicity: str | None = None
grad_school_major: str | None = None
grad_time: str | None = None
@router.get("/api/users", summary="本园区用户列表")
async def park_users(db=Depends(_carrier_db), user: dict = Depends(_carrier_user)):
t = await _my_park(user)
return {"items": await db.users.list_by_park(t["id"])}
@router.post("/api/users", summary="新增本园区用户(全局唯一账号)")
async def park_user_create(body: ParkUserCreateBody, db=Depends(_carrier_db), user: dict = Depends(_carrier_user)):
t = await _my_park(user)
uname = body.username.strip()
if not uname:
raise HTTPException(status_code=400, detail="账号必填")
if await db.users.get_by_username(uname):
raise HTTPException(status_code=400, detail="账号已存在")
if not body.password:
raise HTTPException(status_code=400, detail="请设置初始密码")
u = await db.users.create(uname, body.password, nickname=body.nickname, gender=body.gender,
birthday=body.birthday, phone=body.phone, email=body.email, company=body.company,
id_card=body.id_card, ethnicity=body.ethnicity, grad_school_major=body.grad_school_major,
grad_time=body.grad_time, source="park", role="opc_member", affiliation="park",
park_id=t["id"], park_name=t["name"], account_type="park_staff")
return {"ok": True, "user": u}
@router.put("/api/users/{uid}", summary="更新本园区用户资料")
async def park_user_update(uid: str, body: ParkUserUpdateBody, db=Depends(_carrier_db), user: dict = Depends(_carrier_user)):
t = await _my_park(user)
existing = await db.users.get_by_id(uid)
if existing is None or existing.get("park_id") != t["id"]:
raise HTTPException(status_code=403, detail="仅可管理本园区用户")
u = await db.users.update_profile(uid, body.model_dump(exclude_none=True))
return {"ok": True, "user": u}
# ==================== 园区端:发布资讯 + 大屏内容 ====================
class ParkContentBody(BaseModel):
type: str = "news"
title: str
summary: str = ""
body: str = ""
cover: str = ""
video: str = ""
video_cover: str = ""
link: str = ""
card_mode: str = "small" # 园区默认小图;运营端可改大/小图
@router.post("/api/content", summary="园区端发布资讯(发布人=园区名,默认不公开)")
async def carrier_publish_content(
body: ParkContentBody,
request: Request,
db=Depends(_carrier_db),
user: dict = Depends(_carrier_user),
):
t = await _my_park(user)
item = await db.content.create({
"type": body.type, "title": body.title, "summary": body.summary, "body": body.body,
"cover": body.cover, "video": body.video, "video_cover": body.video_cover, "link": body.link,
"publisher_name": t["name"], "publisher_avatar": "", # 发布人强制=园区名,不可改
"source": "carrier", "tenant_id": t["id"],
"card_mode": body.card_mode or "small",
"is_public": False, # 默认不公开:审核通过后由运营端切公开(进 C 端)
"status": "pending", # 待审核:需运营端在资讯审核通过后才上线
})
await write_audit(db, action="content.publish", resource="content", resource_id=item["id"],
detail=f"carrier {item['title']}", user=user, request=request)
return {"ok": True, "content": item}
@router.get("/api/content", summary="园区资讯列表(本园区,含未公开,供大屏)")
async def carrier_content_list(
db=Depends(_carrier_db),
user: dict = Depends(_carrier_user),
):
t = await _my_park(user)
return {"items": await db.content.list(tenant_id=t["id"])}
@router.get("/api/content/feed", summary="园区资讯流(大屏,tenant 解析,含未公开,直接返回列表)")
async def park_content_feed(
authorization: str | None = Header(default=None),
tenant_id: str | None = None,
):
from app.infrastructure.repositories import Database
tid = await _resolve_tenant(authorization, tenant_id)
db = Database()
try:
items = await db.content.list(status="published", tenant_id=tid)
finally:
await db.close()
return {"items": items}