From 4d4d38fe0fb43ce6ccdc6a58c91d532562b57680 Mon Sep 17 00:00:00 2001 From: Pine Date: Mon, 24 Aug 2026 16:57:27 +0800 Subject: [PATCH] =?UTF-8?q?feat(park):=20=E5=9B=AD=E5=8C=BA=E5=AD=90?= =?UTF-8?q?=E5=BA=94=E7=94=A8=20app/park=20=E9=AA=A8=E6=9E=B6=E8=BF=81?= =?UTF-8?q?=E5=85=A5=20=E2=80=94=20dispatcher=20/park=20=E5=89=8D=E7=BC=80?= =?UTF-8?q?=E5=88=86=E6=B5=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 迁自 park-desktop/backend/app:config(数据目录→serverdata/park)/storage/sim_engine/ mqtt/event_bus + 精简 routers(health/settings/config/dashboard/park/companies/display/ events) + app.py 子应用(lifespan 起 MQTT hub+tick 线程, include_router prefix=/park)。 dispatcher.py 增 ROUTE_PARK_PREFIX + _prefix_app,与 app.training(/api)、app.main 分流。 pyproject 增 paho-mqtt(MQTT 必需)。AI/知识库/语音/视觉重模块延后全量迁入。 --- app/park/__init__.py | 0 app/park/app.py | 67 +++++++++++ app/park/config.py | 106 +++++++++++++++++ app/park/event_bus.py | 44 +++++++ app/park/mqtt.py | 242 +++++++++++++++++++++++++++++++++++++ app/park/routers.py | 179 ++++++++++++++++++++++++++++ app/park/sim_engine.py | 262 +++++++++++++++++++++++++++++++++++++++++ app/park/storage.py | 115 ++++++++++++++++++ dispatcher.py | 22 +++- pyproject.toml | 1 + uv.lock | 8 ++ 11 files changed, 1041 insertions(+), 5 deletions(-) create mode 100644 app/park/__init__.py create mode 100644 app/park/app.py create mode 100644 app/park/config.py create mode 100644 app/park/event_bus.py create mode 100644 app/park/mqtt.py create mode 100644 app/park/routers.py create mode 100644 app/park/sim_engine.py create mode 100644 app/park/storage.py diff --git a/app/park/__init__.py b/app/park/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/app/park/app.py b/app/park/app.py new file mode 100644 index 0000000..a9709fd --- /dev/null +++ b/app/park/app.py @@ -0,0 +1,67 @@ +# -*- coding: utf-8 -*- +"""园区子应用(app/park)—— FastAPI 子应用组装。 + +由 server-core/dispatcher.py 以 `/park` 前缀路由对外;独立 lifespan: +启动 MQTT hub、数据 tick 线程(sim_engine → publish_tick);s2s 实时语音栈延后。 +""" +from __future__ import annotations + +import asyncio +import logging +import threading +from contextlib import asynccontextmanager +from pathlib import Path + +from fastapi import FastAPI +from fastapi.staticfiles import StaticFiles + +from .config import settings +from .event_bus import bus +from .mqtt import hub +from .routers import router +from .sim_engine import sim_engine + +log = logging.getLogger("dpm.park") + +settings.MEDIA_DIR.mkdir(parents=True, exist_ok=True) + + +@asynccontextmanager +async def lifespan(app: FastAPI): + bus.set_loop(asyncio.get_running_loop()) + hub.start() + + stop = threading.Event() + + def tick_loop(): + while not stop.is_set(): + try: + snap = sim_engine.tick() + hub.publish_tick(snap) + except Exception as e: # noqa: BLE001 + log.warning("园区数据 tick 异常: %s", e) + stop.wait(settings.MQTT_TICK_INTERVAL) + + threading.Thread(target=tick_loop, daemon=True).start() + log.info("园区子应用 /park 启动(MQTT %s:%s 系 %s)", + settings.MQTT_HOST, settings.MQTT_PORT, + "已启用" if settings.MQTT_ENABLED else "已禁用") + try: + yield + finally: + stop.set() + hub.stop() + + +app = FastAPI(title="OPC 智能园区子应用", version="1.0.0", lifespan=lifespan) + +# 关键:内部路由本为 /api/*,挂载到 /park 下 → 外部地址 /park/api/* +app.include_router(router, prefix="/park") + +# 媒体静态服务(上传/播放文件) +app.mount("/park/file", StaticFiles(directory=str(settings.MEDIA_DIR)), name="park-media") + +# 管理后台静态资源(如有) +_static_dir = Path(__file__).resolve().parent / "static" +if _static_dir.exists(): + app.mount("/park/static", StaticFiles(directory=str(_static_dir)), name="park-static") diff --git a/app/park/config.py b/app/park/config.py new file mode 100644 index 0000000..7ecb39c --- /dev/null +++ b/app/park/config.py @@ -0,0 +1,106 @@ +# -*- coding: utf-8 -*- +"""园区子应用(app/park)—— 全局配置 + 迁自 park-desktop/backend/app/config.py,数据目录改落 server-core/serverdata/park, + 环境变量沿用 server-core/.env(前缀 DPM_ 或 PARK_ 均兼容,MIG/PARK 二选一)。 +""" + +import os +from pathlib import Path + +SERVER_ROOT = Path(__file__).resolve().parent.parent.parent # server-core 根 +PARK_DIR = Path(os.environ.get("PARK_DATA_DIR", str(SERVER_ROOT / "serverdata" / "park"))) + + +def _load_dotenv(path): + """极简 .env 解析:KEY=VALUE,支持 # 注释与引号""" + if not Path(path).exists(): + return + for raw in Path(path).read_text("utf-8").splitlines(): + line = raw.strip() + if not line or line.startswith("#") or "=" not in line: + continue + key, _, value = line.partition("=") + key = key.strip() + value = value.strip().strip('"').strip("'") + if key and key not in os.environ: + os.environ[key] = value + + +_load_dotenv(SERVER_ROOT / ".env") + + +def _env(key, default): + return os.getenv(key, default) + + +class Settings: + # ---- HTTP ---- + HOST = _env("DPM_HOST", _env("FASTAPI_HOST", "0.0.0.0")) + PORT = int(_env("DPM_PORT", _env("FASTAPI_PORT", "8000"))) + + # ---- 数据目录(统一落 serverdata/park)---- + MEDIA_DIR = Path(_env("DPM_MEDIA_DIR", str(PARK_DIR / "media"))) + DATA_FILE = Path(_env("DPM_DATA_FILE", str(PARK_DIR / "data.json"))) + DIST_DIR = Path(_env("DPM_DIST_DIR", str(PARK_DIR / "dist"))) # 大屏前端产物(可选托管) + + # ---- 视频帧保存(摄像头识别帧)---- + SAVE_VISION_FRAMES = _env("DPM_SAVE_VISION_FRAMES", "1") == "1" + VISION_FRAME_SAVE_DIR = Path(_env("DPM_VISION_SAVE_DIR", str(MEDIA_DIR / "video"))) + + # ---- MQTT(EMQX broker;凭据走 server-core/.env,勿入库)---- + MQTT_ENABLED = _env("DPM_MQTT_ENABLED", "1") == "1" + MQTT_HOST = _env("DPM_MQTT_HOST", _env("MQTT_BROKER_HOST", "192.168.1.9")) + MQTT_PORT = int(_env("DPM_MQTT_PORT", _env("MQTT_BROKER_PORT", "1883"))) + MQTT_USERNAME = _env("DPM_MQTT_USERNAME", _env("MQTT_USERNAME", "")) or None + MQTT_PASSWORD = _env("DPM_MQTT_PASSWORD", _env("MQTT_PASSWORD", "")) or None + MQTT_CLIENT_ID = _env("DPM_MQTT_CLIENT_ID", "dpm-backend") + MQTT_WS_URL = _env( + "DPM_MQTT_WS", + f"ws://{_env('MQTT_PUBLIC_HOST', _env('MQTT_BROKER_HOST', 'localhost'))}:8083/mqtt", + ) + MQTT_TICK_INTERVAL = float(_env("DPM_MQTT_TICK_INTERVAL", "2.2")) + + # ---- 主题 ---- + TOPIC_COMMAND = _env("DPM_TOPIC_COMMAND", "opc/display/command") + TOPIC_TICK = _env("DPM_TOPIC_TICK", "opc/dashboard/tick") + TOPIC_ACK = _env("DPM_TOPIC_ACK", "opc/display/ack") + TOPIC_HEARTBEAT = _env("DPM_TOPIC_HEARTBEAT", "opc/display/heartbeat") + SCREEN_TTL = float(_env("DPM_SCREEN_TTL", "30")) + + # ---- 阿里云 DashScope(LLM + 语音识别)---- + DASHSCOPE_API_KEY = _env("DASHSCOPE_API_KEY", "") + LLM_MODEL = _env("DPM_LLM_MODEL", "qwen-plus") + ASR_MODEL = _env("DPM_ASR_MODEL", "paraformer-realtime-v2") + LLM_BASE_URL = _env("DPM_LLM_BASE_URL", "https://dashscope.aliyuncs.com/compatible-mode/v1") + ASR_TIMEOUT = float(_env("DPM_ASR_TIMEOUT", "30")) + + # ---- 管理后台 ---- + ADMIN_SECRET = _env("DPM_ADMIN_SECRET", "dpm-admin-secret-change-me") + + # ---- s2s 实时语音栈(VAD本地 + 云ASR/LLM/TTS)---- + S2S_ENABLED = _env("DPM_S2S_ENABLED", "0") == "1" + S2S_HOST = _env("DPM_S2S_HOST", "0.0.0.0") + S2S_PORT = int(_env("DPM_S2S_PORT", "8765")) + S2S_NUM_PIPELINES = int(_env("DPM_S2S_NUM_PIPELINES", "1")) + VISION_LLM_MODEL = _env("DPM_VISION_LLM_MODEL", "qwen3-vl-flash") + S2S_STT_MODEL = _env("DPM_S2S_STT_MODEL", "qwen3-asr-flash-realtime") + S2S_LLM_MODEL = _env("DPM_S2S_LLM_MODEL", "qwen-plus") + S2S_TTS_MODEL = _env("DPM_S2S_TTS_MODEL", "qwen3-tts-flash-realtime") + S2S_TTS_VOICE = _env("DPM_S2S_TTS_VOICE", "Cherry") + S2S_BUILD_RETRIES = int(_env("DPM_S2S_BUILD_RETRIES", "4")) + S2S_BUILD_RETRY_DELAY = float(_env("DPM_S2S_BUILD_RETRY_DELAY", "10")) + S2S_LOG_LEVEL = _env("DPM_S2S_LOG_LEVEL", "INFO").upper() + VOICE_WS = _env("DPM_VOICE_WS", "") + S2S_WS_URL = ( + VOICE_WS + or f"ws://{_env('MQTT_PUBLIC_HOST', _env('MQTT_BROKER_HOST', 'localhost'))}:{S2S_PORT}/v1/realtime" + ) + S2S_INSTRUCTIONS = _env( + "DPM_S2S_INSTRUCTIONS", + "你是昆小创,昆明市大学生创业园的专属智能语音助手。可用工具:get_park_overview(园区实时数据)、" + "query_companies(企业名录)、control_display(大屏控制)、get_time(时间)。" + "涉及园区数据/企业/大屏控制时务必调用工具获取准确信息。请用简洁专业的中文回答,不超过三句话。", + ) + + +settings = Settings() diff --git a/app/park/event_bus.py b/app/park/event_bus.py new file mode 100644 index 0000000..6b0f054 --- /dev/null +++ b/app/park/event_bus.py @@ -0,0 +1,44 @@ +# -*- coding: utf-8 -*- +"""SSE 事件总线 —— 兼容旧通道(前端 MQTT 不可用时回退) + 所有 MQTT 控制指令同时推送到 SSE /api/events +""" + +import asyncio +import json + + +class EventBus: + def __init__(self): + self._queues = set() + self._loop = None + + def set_loop(self, loop): + self._loop = loop + + def subscribe(self): + q = asyncio.Queue(maxsize=256) + self._queues.add(q) + return q + + def unsubscribe(self, q): + self._queues.discard(q) + + def emit(self, payload: dict): + """线程安全:从任意线程广播事件""" + if not self._loop: + return + try: + asyncio.run_coroutine_threadsafe(self._emit(payload), self._loop) + except Exception: # noqa: BLE001 + pass + + async def _emit(self, payload): + data = json.dumps(payload, ensure_ascii=False) + for q in list(self._queues): + try: + q.put_nowait(data) + except asyncio.QueueFull: + pass + + +bus = EventBus() diff --git a/app/park/mqtt.py b/app/park/mqtt.py new file mode 100644 index 0000000..81c5690 --- /dev/null +++ b/app/park/mqtt.py @@ -0,0 +1,242 @@ +# -*- coding: utf-8 -*- +"""MQTT 发布/订阅中心 —— 后端 ↔ 前端控制通道 + 发布:opc/display/command(页面/媒体/卡片控制)、opc/dashboard/tick(数据快照) + 订阅:opc/display/heartbeat(大屏在线心跳,统计在线大屏数) +""" + +import json +import logging +import os +import threading +import time +import uuid + +import paho.mqtt.client as mqtt + +from .config import settings +from .event_bus import bus + +log = logging.getLogger("dpm.mqtt") + + +def new_cmd_id(): + return uuid.uuid4().hex[:12] + + +def _unique_client_id(base): + """每个客户端进程使用唯一 client_id,避免与其它实例竞争被 broker 踢下线""" + return f"{base}-{os.getpid()}-{uuid.uuid4().hex[:6]}" + + +class MqttHub: + def __init__(self): + self.client = None + self.connected = False + self._lock = threading.RLock() + self.devices = {} # device_id(client_id) -> {role, parent_id, last_seen} + self.last_command = None # {action, params, ts, published} + + # ---------- 生命周期 ---------- + def start(self): + if not settings.MQTT_ENABLED: + log.info("MQTT 已禁用(DPM_MQTT_ENABLED=0)") + return + try: + self.client = mqtt.Client( + client_id=_unique_client_id(settings.MQTT_CLIENT_ID), + clean_session=True, + protocol=mqtt.MQTTv311, + ) + if settings.MQTT_USERNAME: + self.client.username_pw_set(settings.MQTT_USERNAME, settings.MQTT_PASSWORD) + self.client.on_connect = self._on_connect + self.client.on_disconnect = self._on_disconnect + self.client.on_message = self._on_message + self.client.connect_async(settings.MQTT_HOST, settings.MQTT_PORT, keepalive=30) + self.client.loop_start() + log.info("MQTT 连接中 %s:%s ...", settings.MQTT_HOST, settings.MQTT_PORT) + except Exception as e: # noqa: BLE001 + log.warning("MQTT 初始化失败(数据 API 不受影响): %s", e) + + def stop(self): + if self.client: + try: + self.client.loop_stop() + self.client.disconnect() + except Exception: # noqa: BLE001 + pass + self.client = None + + def _on_connect(self, client, userdata, flags, rc): + if rc == 0: + self.connected = True + log.info("MQTT 已连接 %s:%s", settings.MQTT_HOST, settings.MQTT_PORT) + client.subscribe(settings.TOPIC_HEARTBEAT, qos=0) + else: + log.warning("MQTT 连接失败 rc=%s", rc) + + def _on_disconnect(self, client, userdata, rc): + self.connected = False + if rc != 0: + log.warning("MQTT 断开(rc=%s),自动重连中...", rc) + + def _on_message(self, client, userdata, msg): + """接收大屏心跳:更新设备表在线时间(设备由注册表登记 role/parent)""" + if msg.topic == settings.TOPIC_HEARTBEAT: + try: + payload = json.loads(msg.payload.decode("utf-8")) + cid = payload.get("client_id") or msg.topic + with self._lock: + cur = self.devices.get(cid, {"role": "", "parent_id": ""}) + cur["last_seen"] = time.time() + self.devices[cid] = cur + self._prune() + except Exception: # noqa: BLE001 + pass + + def _prune(self): + """剔除超时未上报的设备""" + cutoff = time.time() - settings.SCREEN_TTL + self.devices = {k: v for k, v in self.devices.items() if v.get("last_seen", 0) > cutoff} + + # ---------- 状态查询 ---------- + def screens_online(self): + with self._lock: + return sum(1 for v in self.devices.values() if v.get("last_seen", 0) > time.time() - settings.SCREEN_TTL) + + def status(self): + with self._lock: + self._prune() + devices = [ + { + "device_id": cid, + "role": v.get("role", ""), + "parent_id": v.get("parent_id", ""), + } + for cid, v in self.devices.items() + ] + return { + "mqtt_connected": self.connected, + "mqtt_host": f"{settings.MQTT_HOST}:{settings.MQTT_PORT}", + "screens_online": len(devices), + "devices": devices, + "last_command": self.last_command, + } + + # ---------- 发布 ---------- + def publish(self, topic, payload, qos=1, retain=False): + if not self.client or not self.connected: + log.debug("MQTT 未连接,丢弃发布 %s", topic) + return False + try: + info = self.client.publish(topic, json.dumps(payload, ensure_ascii=False), qos=qos, retain=retain) + return info.rc == mqtt.MQTT_ERR_SUCCESS + except Exception as e: # noqa: BLE001 + log.warning("MQTT 发布失败: %s", e) + return False + + def register_device(self, device_id, role="", parent_id=""): + """前端上报设备:唯一 device_id + 角色(main/secondary) + 从属主屏 parent_id。""" + device_id = (device_id or "").strip() + if not device_id: + return + with self._lock: + cur = self.devices.get(device_id, {}) + if role: + cur["role"] = role + if parent_id: + cur["parent_id"] = parent_id + cur["last_seen"] = time.time() + self.devices[device_id] = cur + self._prune() + log.info("mqtt: 设备注册 id=%s role=%s parent=%s", device_id, role or cur.get("role"), parent_id or cur.get("parent_id")) + + def publish_command(self, action, params=None, screen_id=None, screen_role=None): + """按客户端(设备)路由下发指令: + · screen_id 精确到某台设备 → 只发到该设备独有 topic + · screen_role(main/secondary) → 解析到该角色的设备,发到其独有 topic + · 两者皆空(全部)→ 遍历在线设备列表,依次下发到各自独有 topic + 客户端只订阅自己的独有频道,因此定向指令绝不会被其它屏幕收到。""" + payload = { + "cmd_id": new_cmd_id(), + "ts": int(time.time() * 1000), + "action": action, + "params": params or {}, + } + target = None + requested = bool(screen_id or screen_role in ("main", "secondary")) + if screen_id: + target = (screen_id or "").strip() + elif screen_role in ("main", "secondary"): + with self._lock: + for cid, v in self.devices.items(): + if v.get("role") == screen_role: + target = cid + break + if target: + # 定向:只发到目标设备独有频道 + return self._publish_to(payload, target, screen_role) + if requested: + # 指定了目标但未找到对应设备:不广播,避免误发到全部 + log.warning("mqtt: 目标设备未找到(id=%s role=%s),未下发", screen_id, screen_role) + payload["published"] = False + with self._lock: + self.last_command = {"action": action, "params": params or {}, "ts": payload["ts"], "published": False, "topic": "none"} + return payload + + # 全部:遍历在线设备列表,依次下发到各自独有频道 + with self._lock: + cutoff = time.time() - settings.SCREEN_TTL + online = [cid for cid, v in self.devices.items() if v.get("last_seen", 0) > cutoff] + ok = True + published_topics = [] + for cid in online: + p = dict(payload) + p["screen_id"] = cid + with self._lock: + p["screen_role"] = self.devices.get(cid, {}).get("role", "") + if self.publish(f"{settings.TOPIC_COMMAND}/{cid}", p): + bus.emit(p) + published_topics.append(cid) + else: + ok = False + payload["published"] = ok + payload["screen_id"] = "" + payload["screen_role"] = "both" + with self._lock: + self.last_command = { + "action": action, + "params": params or {}, + "ts": payload["ts"], + "published": ok, + "topic": "all", + } + log.info("mqtt: 全部下发 %s 到 %d 台设备", action, len(published_topics)) + return payload + + def _publish_to(self, payload, target, screen_role=""): + """发到指定设备独有频道,并记录 last_command。""" + topic = f"{settings.TOPIC_COMMAND}/{target}" + payload["screen_id"] = target + with self._lock: + payload["screen_role"] = self.devices.get(target, {}).get("role", screen_role or "") + ok = self.publish(topic, payload) + payload["published"] = ok + with self._lock: + self.last_command = { + "action": payload["action"], + "params": payload.get("params") or {}, + "ts": payload["ts"], + "published": ok, + "topic": topic, + } + if ok: + bus.emit(payload) + return payload + + def publish_tick(self, snapshot): + payload = {"ts": int(time.time() * 1000), "snapshot": snapshot} + self.publish(settings.TOPIC_TICK, payload, qos=0) + + +hub = MqttHub() diff --git a/app/park/routers.py b/app/park/routers.py new file mode 100644 index 0000000..61466a3 --- /dev/null +++ b/app/park/routers.py @@ -0,0 +1,179 @@ +# -*- 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 time + +from fastapi import APIRouter, Request +from fastapi.responses import JSONResponse, StreamingResponse +from pydantic import BaseModel, Field + +from .config import settings +from .event_bus import bus +from .mqtt import hub +from .sim_engine import sim_engine +from .storage import storage + +log = logging.getLogger("dpm.api") + +router = APIRouter() + + +# ==================== 数据模型 ==================== + +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 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 = "" + + +# ==================== 健康 / 设置 / 引导 ==================== + +@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", + } + + +@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/语音接入地址与 tick 间隔。""" + return { + "mqtt": { + "host": settings.MQTT_HOST, + "port": settings.MQTT_PORT, + "ws_url": settings.MQTT_WS_URL, + "tick_interval": settings.MQTT_TICK_INTERVAL, + }, + "voice_ws": settings.S2S_WS_URL if settings.S2S_ENABLED else "", + "media_base": f"{request.base_url}park/file", + "api_base": f"{request.base_url}park/api", + } + + +# ==================== 大屏数据(sim_engine) ==================== + +@router.get("/api/dashboard/snapshot") +async def dashboard_snapshot(): + return sim_engine.snapshot() + + +@router.get("/api/dashboard/overview") +async def dashboard_overview(): + return sim_engine.snapshot() + + +@router.get("/api/park/companies") +async def park_companies(): + snap = sim_engine.snapshot() + return snap.get("companies", []) + + +@router.get("/api/park/zones") +async def park_zones(): + snap = sim_engine.snapshot() + return snap.get("zones", []) + + +# ==================== 展示控制(管理端 → 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): + hub.register_device(body.device_id, body.role, body.parent_id) + return {"ok": True, "device_id": body.device_id, "role": body.role, "parent_id": body.parent_id} + + +@router.post("/api/display/command") +async def display_command_publish(body: DisplayCommandBody): + 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) + 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", + "mqtt_connected": hub.connected} + + +@router.get("/api/display/state") +async def display_state(): + return hub.status() + + +# ==================== 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.timeout(3.11+ 的上下文管理器)。""" + import asyncio + return asyncio.wait_for(awaitable, timeout=seconds) diff --git a/app/park/sim_engine.py b/app/park/sim_engine.py new file mode 100644 index 0000000..d777048 --- /dev/null +++ b/app/park/sim_engine.py @@ -0,0 +1,262 @@ +# -*- coding: utf-8 -*- +"""园区数据引擎(后端)—— 真实数据 + 测算演示 + 数据来源(用户指定口径): + · 《昆明市大学生创业园运行情况统计表(2026年7月).xls》(官方月报) + · 《昆明市创业园宣传册.pdf》 + 测算口径(用户确认): + · TOKEN:真实业务锚点月均 200 亿 tokens → 日均 ≈ 6.67 亿(200亿/30), + 实时速率 ≈ 7,700 t/s(日均/86400s),近30日累计 = 200 亿 + · 工具调用:按平均每次调用消耗约 2 万 tokens 折算(今日 ≈ 3.3 万次) + 其余统计指标为真实静态值;TOKEN/工具仅做轻微演示波动 +""" + +import random +import threading +import time + +# ---------- 真实锚点(仅来自 xls 统计表 + 附件1 省级材料 + 宣传册) ---------- +REAL = { + "as_of": "2026-07", + "token_monthly_yi": 200, # 月均消耗 200 亿 tokens(AI 平台业务锚点,非政府统计口径) + "token_daily_wan": 66666.7, # 日均 6.67 亿 tokens(200亿/30 ≈ 66666.7 万) + "token_rate": 7700, # 实时 t/s(6.67e8 / 86400 ≈ 7,716) + "tools_per_call_tokens": 150000, # 平均每次工具调用消耗 tokens(真实测算:12.5M输入 ÷ 81步 ≈ 15.4万) + "occupied": 38, # 实际入驻企业数(官方口径:《2026年7月运行情况统计表》38 家) + "capacity": 49, # 园区可容纳企业个数(统计表) + "invested": 560, # 开园以来累计投入运营资金(万元,统计表) + "jobs": 213, # 带动就业人数(统计表,当年累计口径) + "revenue_total": 252.15, # 生产经营总额(万元,统计表,当年累计口径) + "revenue_tax": 1.03, # 上缴税利总额(万元,统计表,当年累计口径) + "incubated_total": 181, # 累计孵化企业(附件1 省级绩效累计口径,截至2026年5月) + "in_incubation": 31, # 当前实有在孵创业实体(附件1) + "graduated": 149, # 累计成功孵化出园企业(统计表,开园以来累计) + "patents": 171, # 入驻团队发明专利(附件1 累计口径) + "jobs_total": 1529, # 累计带动就业(附件1 省级绩效累计口径) + "revenue_total_cum": 21911.14, # 入驻团队累计经营收入(万元,附件1 累计口径 ≈ 2.19 亿元) + "revenue_tax_cum": 426.2, # 入驻团队累计税收(万元,附件1 累计口径) + "mentors": 50, # 创业指导专家团队(附件1) + "success_rate": 95.91, # 近3年入孵实体孵化成功率(%)(附件1) + "area": 3000, # 园区建筑面积(㎡,统计表+宣传册) + "founded": 2009, # 成立年份(宣传册) + "province_level": 2012, # 获评省级创业园年份(宣传册) + "address": "昆明市民航路229号(人力资源中心1-2楼)", + "phone": "18087174891", + "email": "714066050@qq.com", +} + +# 入驻企业名录(《入驻企业信息表》39 家全称,与在园数统一口径) +COMPANIES = [ + "云南派音人工智能科技有限公司", "米勒克尔蓝宝石珠宝产业化(项目)", "中泰研学合作(项目)", + "云南宸中低空经济有限公司", "昆明智海银高文化科技有限公司", "人工智能机器人大模型训练(项目)", + "云南廷秀文旅康养服务(项目)", "瀚颖AI+教育信息咨询(项目)", "云南大学AI+联合创业服务平台(项目)", + "仰光客厅(项目)", "云南上古绝学文化发展有限公司", "千楣之路(项目)", + "中越生物医疗(项目)", "中越文旅交流(项目)", "绿态环保科技项目", + "联翼航空技术低空经济应用国际化拓展(项目)", "酷享野农AI农业(项目)", "滇缅国际设计(项目)", + "昆明舒诺生物科技有限责任公司", "达岸教育管理(云南)有限公司", "昆明滇油石化有限责任公司", + "花仙子园艺肥料(云南)有限公司", "Facebook越南市场普洱茶与建水紫陶壶跨境电商(项目)", + "研X——同行者网络(项目)", "鬼才明AI创意工作室(项目)", "朵哈·玫瑰特色产业链(项目)", + "昆明云韵体育有限责任公司", "昆明舒护安养老服务有限责任公司", + "南菌优培—食用菌专用生物制剂创新与高原特色野生菌高值化全链开发(项目)", + "五华区丽裳文化艺术工作室", "云南星瑞航空(项目)", "综合性云南直播私域服务平台", + "昆明屿澈电子商务有限公司", "中外青少年研学(项目)", "星杭商务信息咨询昆明有限责任公司", + "蓝智科技(项目)", "云品出滇·纸享万家(项目)", + "云迹轻旅—昆明本土一站式轻量化文旅服务工作室(项目)", "启元人工智能科技(项目)", +] + +# 宣传册四大孵化区域 +ZONES = [ + {"name": "OPC创业空间", "desc": "AI 技术支撑 · 数字经济/AI应用/软件开发/文创设计/电商直播等轻资产领域", "period": "拎包入驻 · 最长2年"}, + {"name": "创业苗圃区", "desc": "创意阶段 · 未注册企业团队 · 开放式共享工位", "period": "最长2年"}, + {"name": "孵化加速区", "desc": "已工商注册初创企业 · 独立办公空间", "period": "基础2年 · 可延1年(≤3年)"}, + {"name": "国际创客区", "desc": "归国留学生 · 外籍留学生 · 国际化配套服务", "period": "参照加速区"}, +] + +# 园区企业行业分布(按《入驻企业信息表》39 家企业简介分类统计) +INDUSTRY_MIX = [ + {"name": "人工智能", "value": 7}, + {"name": "跨境电商", "value": 5}, + {"name": "教育培训", "value": 4}, + {"name": "文旅康养", "value": 4}, + {"name": "文化创意", "value": 4}, + {"name": "生物医药", "value": 4}, + {"name": "低空经济", "value": 3}, + {"name": "现代农业", "value": 3}, + {"name": "绿色环保", "value": 2}, + {"name": "企业服务", "value": 2}, + {"name": "新材料", "value": 1}, +] + +# 园区实时动态(真实事件:园区官方公示 + 宣传册活动 + 互联网公开信息,最新在前) +FEED = [ + {"id": 24, "icon": "dna", "text": "园内企业「昆明舒诺生物科技」完成GEO服务体系建设(生成式引擎优化,聚焦宠物营养品牌数字化)", "time": "2026-08"}, + {"id": 1, "icon": "activity", "text": "2026年7月运行情况统计表填报完成:累计投入运营资金560万元、带动就业213人", "time": "2026-07-31"}, + {"id": 2, "icon": "radio", "text": "2026年第三期入驻创业企业(项目)评审结果在昆明市人社局官网公示", "time": "2026-07-02"}, + {"id": 3, "icon": "graduation", "text": "云南师范大学就业创业实践基地落地园区 · 云南旅游职业学院开展就业创业交流", "time": "2026-06-03"}, + {"id": 23, "icon": "radio", "text": "「PineSound」完成 ICP 备案(滇ICP备2026008517号)与公安备案,注册地址昆明市民航路301-307号", "time": "2026"}, + {"id": 4, "icon": "rocket", "text": "\"滇水逐梦·雨林创航\"昆明·西双版纳创业路演PK赛成功举办", "time": "2026-05-20"}, + {"id": 5, "icon": "flag", "text": "昆明市青年企业家协会创新创业基地在园区揭牌", "time": "2026-05-13"}, + {"id": 6, "icon": "users", "text": "\"聚力赋能·共创未来\"高校创新创业交流暨项目路演活动圆满举行", "time": "2026-05-08"}, + {"id": 7, "icon": "trophy", "text": "\"创赢未来\"2026创业大赛昆明选拔赛暨马兰花创业培训讲师大赛举行(13个项目参赛)", "time": "2026-04-22"}, + {"id": 8, "icon": "broadcast", "text": "2026年第二期创业企业(项目)入驻招募公告发布(官网 www.kmyc.gov.cn)", "time": "2026-04-15"}, + {"id": 21, "icon": "sparkles", "text": "「PineSound」AI驱动音频管理平台上线 pinesound.cn:智能配乐与音效生成,集成100W+音效库、50W+配乐库,全球版权授权(官网)", "time": "2026"}, + {"id": 22, "icon": "medal", "text": "PineSound发布企业标准 Q/YNPY 001-2026《数字内容创作 配乐通用分类标准》", "time": "2026"}, + {"id": 9, "icon": "award", "text": "2026年第一期入驻评审结果在昆明市人社局官网公示(23个优质项目正式入驻)", "time": "2026-03-27"}, + {"id": 10, "icon": "users", "text": "2026年第一期招募集中评审完成:38个申请,23个优质项目正式入驻", "time": "2026-03-24"}, + {"id": 19, "icon": "sparkles", "text": "云南派音人工智能科技专注多模态音频技术研发:自研Pine系列模型覆盖音频识别、向量嵌入、音效生成、配乐创作", "time": "2026-04"}, + {"id": 20, "icon": "bolt", "text": "「PineSound」2026年4月注册成立,获北京投资支持并吸纳就业4人,入驻加速区(入驻企业信息表)", "time": "2026-04"}, + {"id": 11, "icon": "cpu", "text": "云南首个人工智能OPC创新人才基地落地昆明,填补省内个体AI创业培育空白(公开报道)", "time": "2026-03-23"}, + {"id": 12, "icon": "medal", "text": "第九届\"春城创业荟\"创业创新大赛圆满闭幕,获奖项目名单公布", "time": "2025-09-30"}, + {"id": 13, "icon": "trophy", "text": "\"春城创业荟\"初赛落幕,147个项目晋级复赛", "time": "2025-04-28"}, + {"id": 14, "icon": "building", "text": "园内企业「花仙子园艺肥料(云南)」成立,注册资本118万元(公开报道)", "time": "2025-04-28"}, + {"id": 15, "icon": "radio", "text": "2025年入园项目评审结果公示(昆明市人社局官网)", "time": "2025-04-22"}, + {"id": 16, "icon": "atom", "text": "研X平台入选\"创客中国\"项目库(一站式研究生学术成长平台)", "time": "2025-01"}, + {"id": 17, "icon": "leaf", "text": "行业动态:园内「朵哈·玫瑰」原料产地昆明八街食用玫瑰行情上涨,玫瑰经济走强", "time": "2025-01"}, + {"id": 18, "icon": "heart", "text": "行业动态:「千楣之路」主营檀娜卡——泰国天然护肤品牌Thanaka正式进入中国市场", "time": "2025-01"}, +] + + +def _day_secs(): + """当日零点起已过秒数(时钟驱动,保证数值单向持续增长、重启连续)""" + now = time.time() + lt = time.localtime(now) + day_start = time.mktime((lt.tm_year, lt.tm_mon, lt.tm_mday, 0, 0, 0, 0, 0, -1)) + return max(0.0, now - day_start) + + +def _today_values(prev_noise=None): + """按真实速率推算当日累计值(单向增长 + 工具调用波动): + 平台实时速率 = 单会话 157 tok/s × 49 路并发 ≈ 7,700 t/s + 今日 TOKEN(万)= 秒数 × 7,700 / 10000 = 秒数 × 0.77 + 今日工具调用 = 今日 TOKEN / 15万 tokens 每次(真实基线) + + 随机游走噪声(±10,波浪式波动,不单调) + """ + secs = _day_secs() + today_wan = round(secs * 0.77) # 万 tokens + base_tools = secs * 0.77 * 10000 / REAL["tools_per_call_tokens"] # 真实基线(浮点) + if prev_noise is None: + noise = 0.0 + else: + noise = max(-10.0, min(10.0, prev_noise + (random.random() - 0.5) * 6.0)) + today_tools = max(0, round(base_tools + noise)) # 次 + return today_wan, today_tools, noise + + +def _make_series(base, noise, n=30): + """生成围绕日均锚点的演示曲线(单位:万 tokens),尾点=今日实时值""" + out = [] + for _ in range(n - 1): + v = base + (random.random() - 0.5) * noise + out.append(round(max(1, v), 1)) + out.append(round(base, 1)) + return out + + +def _live_values(prev=None): + """TOKEN 面板实时运行指标(围绕真实基准动态波动,会话持续推进)""" + r = random.random + steps = (prev or {}).get("steps", 81) + steps += 1 if r() < 0.3 else 0 + return { + "firstToken": round(1.3 + (r() - 0.5) * 0.2, 2), # 首 token 平均 1.3s 基准 ±0.1 + "throughput": round(157 + (r() - 0.5) * 14), # 吞吐 157 tok/s 基准 ±7 + "cacheHit": round(98 + (r() - 0.5) * 1.2, 1), # 缓存命中 98% 基准 ±0.6 + "inputM": round((prev or {}).get("inputM", 12.5) + 0.004 + r() * 0.008, 2), # 会话输入持续推进 + "outputM": round((prev or {}).get("outputM", 2.0) + 0.001 + r() * 0.002, 2), # 会话输出持续推进 + "rounds": max(2, steps // 40), # 每 40 步一轮 + "steps": steps, # 步数持续增长 + } + + +def init_snapshot(): + today_wan, today_tools, noise = _today_values() + return { + "as_of": REAL["as_of"], + "_toolsNoise": noise, + "live": _live_values(), + "token": { + "today": today_wan, # 万 tokens(当日累计,实时增长) + "total": REAL["token_monthly_yi"] * 10000 + today_wan, # 万 tokens(近30日 = 200亿 + 今日增量) + "rate": REAL["token_rate"], # t/s(157 tok/s × 49 并发 ≈ 7,700) + "series": _make_series(REAL["token_daily_wan"], REAL["token_daily_wan"] * 0.16), + }, + "tools": { + "today": today_tools, # 次(随 TOKEN 实时联动) + "total": round(REAL["token_monthly_yi"] * 100000000 / REAL["tools_per_call_tokens"]) + today_tools, # 次 + "success": 98.5, # 成功率(演示) + "nodes": 12, # AI 服务节点(演示) + }, + "projects": { + "inPark": REAL["occupied"], # 实际入驻 38 家(2026年7月统计表官方口径) + "capacity": REAL["capacity"], # 可容纳 49 个 + "invested": REAL["invested"], # 累计投入 560 万元 + }, + "jobs": {"total": REAL["jobs"]}, # 带动就业 213 人(当年累计口径) + "revenue": {"total": REAL["revenue_total"], "tax": REAL["revenue_tax"]}, # 万元(当年累计口径) + "cumulative": { # 省级绩效累计口径(附件1,截至2026年5月) + "incubated": REAL["incubated_total"], # 累计孵化企业 181 家 + "inIncubation": REAL["in_incubation"], # 当前实有在孵 31 家 + "graduated": REAL["graduated"], # 累计成功出园 149 家 + "patents": REAL["patents"], # 发明专利 171 项 + "jobs": REAL["jobs_total"], # 累计带动就业 1,529 人 + "revenue": REAL["revenue_total_cum"], # 累计经营收入 21,911.14 万元 + "tax": REAL["revenue_tax_cum"], # 累计税收 426.2 万元 + "mentors": REAL["mentors"], # 创业指导专家团队 50 人 + "successRate": REAL["success_rate"], # 近3年孵化成功率 95.91% + }, + "park": { + "area": REAL["area"], + "founded": REAL["founded"], + "provinceLevel": REAL["province_level"], + "address": REAL["address"], + "phone": REAL["phone"], + "email": REAL["email"], + }, + "zones": ZONES, + "companies": COMPANIES, + "industryMix": INDUSTRY_MIX, + "feed": FEED, + } + + +def next_snapshot(s): + """时钟驱动:按真实速率重算当日累计(TOKEN 单向增长,工具调用波浪波动),其余真实静态""" + today_wan, today_tools, noise = _today_values(s.get("_toolsNoise")) + token = { + "today": today_wan, + "total": REAL["token_monthly_yi"] * 10000 + today_wan, + "rate": REAL["token_rate"], + "series": s["token"]["series"][:-1] + [today_wan], + } + tools = { + "today": today_tools, + "total": round(REAL["token_monthly_yi"] * 100000000 / REAL["tools_per_call_tokens"]) + today_tools, + "success": 98.5, + "nodes": 12, + } + return { + **s, + "_toolsNoise": noise, + "live": _live_values(s.get("live")), + "token": token, + "tools": tools, + } + + +class SimEngine: + """带锁的快照引擎:单例供 API 与 MQTT tick 共用""" + + def __init__(self): + self._lock = threading.RLock() + self._snap = init_snapshot() + + def snapshot(self): + with self._lock: + return dict(self._snap) + + def tick(self): + with self._lock: + self._snap = next_snapshot(self._snap) + return dict(self._snap) + + +sim_engine = SimEngine() diff --git a/app/park/storage.py b/app/park/storage.py new file mode 100644 index 0000000..28f3a8d --- /dev/null +++ b/app/park/storage.py @@ -0,0 +1,115 @@ +# -*- coding: utf-8 -*- +"""JSON 持久化:设置 / 播放列表 / URL 媒体(对齐原 Rust storage 语义)""" + +import json +import threading +import time +from pathlib import Path + +from .config import settings + +DEFAULT_SETTINGS = { + "username": "admin", + "password": "123456", + "volume": 80, + "sfx_volume": 60, + "auto_play": True, + "play_mode": "sequential", + "image_duration": 5, + "fullscreen": True, + "autostart": False, +} + + +class Storage: + def __init__(self, data_file: Path): + self.data_file = data_file + self._lock = threading.RLock() + self._data = self._load() + + # ---------- 底层 ---------- + def _load(self): + if self.data_file.exists(): + try: + data = json.loads(self.data_file.read_text("utf-8")) + # 规范化:URL 来源自动标记(兼容旧数据 source=None 的 URL 条目) + for it in data.setdefault("playlist", []): + p = it.get("path", "") + if p.startswith(("http://", "https://")) and not it.get("source"): + it["source"] = "url" + return data + except Exception: + pass + return {"settings": dict(DEFAULT_SETTINGS), "playlist": [], "url_media": []} + + def _save(self): + self.data_file.parent.mkdir(parents=True, exist_ok=True) + tmp = self.data_file.with_suffix(".tmp") + tmp.write_text(json.dumps(self._data, ensure_ascii=False, indent=2), "utf-8") + tmp.replace(self.data_file) + + # ---------- 设置 ---------- + def get_settings(self): + with self._lock: + s = dict(DEFAULT_SETTINGS) + s.update(self._data.get("settings", {})) + return s + + def update_settings(self, **kw): + with self._lock: + s = dict(DEFAULT_SETTINGS) + s.update(self._data.get("settings", {})) + for k, v in kw.items(): + if v is not None and k in s: + s[k] = v + self._data["settings"] = s + self._save() + + # ---------- 播放列表 ---------- + def get_playlist(self): + with self._lock: + return list(self._data.get("playlist", [])) + + def add_to_playlist(self, path): + with self._lock: + items = self._data.setdefault("playlist", []) + if not any(it.get("path") == path for it in items): + # URL 来源自动标记,避免被当作本地文件过滤 + source = "url" if path.startswith(("http://", "https://")) else None + items.append({"path": path, "source": source, "name": None}) + self._save() + + def remove_from_playlist(self, path): + with self._lock: + self._data.setdefault("playlist", []) + self._data["playlist"] = [it for it in self._data["playlist"] if it.get("path") != path] + self._save() + + # ---------- URL 媒体 ---------- + def get_url_media(self): + with self._lock: + return list(self._data.get("url_media", [])) + + def add_url_media(self, url, name, media_type): + """返回 True 表示新增,False 表示已存在(重复)""" + with self._lock: + items = self._data.setdefault("url_media", []) + if any(it.get("url") == url for it in items): + return False + items.append({ + "url": url, "name": name, "type": media_type, + "added_at": time.strftime("%Y-%m-%d %H:%M:%S"), + }) + self._save() + return True + + def delete_media(self, path): + with self._lock: + self._data.setdefault("url_media", []) + self._data["url_media"] = [it for it in self._data["url_media"] if it.get("url") != path] + self._data.setdefault("playlist", []) + self._data["playlist"] = [it for it in self._data["playlist"] if it.get("path") != path] + self._save() + + +storage = Storage(settings.DATA_FILE) diff --git a/dispatcher.py b/dispatcher.py index 507a3c9..9d716c8 100644 --- a/dispatcher.py +++ b/dispatcher.py @@ -1,8 +1,9 @@ """server-core 统一入口(dispatcher)—— 对外生产服务 opc.pinesound.cn 单入口。 -同一端口服务平台应用与培训子应用,按路径路由: - /api/*、/uploads/*、/SpXvScDiDT.txt → 培训子应用(app.training,保留 /api 路由与独立 data/opc.db) - 其余(/auth /opc /admin /agents 等)→ 平台应用(app.main,身份 + 业务域) +按「端口专属前缀」把请求分流到对应子应用,每端口自带前缀子路由、互不冲突: + /park/* → 园区子应用(app.park,大屏/MQTT/园区数据) + /api/*、/uploads/*、/SpXvScDiDT.txt → 培训子应用(app.training,独立 data/opc.db) + /auth /opc /admin /government /investor …→ 平台应用(app.main,身份 + 业务域,各端口子路由) 启动(server-core 目录): uv run uvicorn dispatcher:app --host 0.0.0.0 --port 8090 --reload @@ -12,11 +13,22 @@ from __future__ import annotations import logging from app.main import app as core_app +from app.park.app import app as park_app from app.training.main import app as training_app logger = logging.getLogger(__name__) ROUTE_TRAINING_PREFIXES = ("/api", "/uploads", "/SpXvScDiDT.txt") +ROUTE_PARK_PREFIX = "/park" + + +def _prefix_app(path: str): + """按路径前缀匹配目标子应用;未命中返回 None 交由平台应用(app.main)。""" + if path.startswith(ROUTE_PARK_PREFIX): + return park_app + if path.startswith(ROUTE_TRAINING_PREFIXES): + return training_app + return None class Dispatcher: @@ -56,8 +68,8 @@ class Dispatcher: await self._run_lifespan(receive, send) return path = scope.get("path", "") - target = training_app if path.startswith(ROUTE_TRAINING_PREFIXES) else core_app + target = _prefix_app(path) or core_app await target(scope, receive, send) -app = Dispatcher([core_app, training_app]) +app = Dispatcher([core_app, training_app, park_app]) diff --git a/pyproject.toml b/pyproject.toml index 2c42b4a..3d24e41 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -16,6 +16,7 @@ dependencies = [ "pyjwt>=2.8.0", "httpx>=0.27.0", "python-multipart>=0.0.12", + "paho-mqtt>=1.6,<2", ] [dependency-groups] diff --git a/uv.lock b/uv.lock index 5d84de8..b2dc677 100644 --- a/uv.lock +++ b/uv.lock @@ -799,6 +799,12 @@ wheels = [ { url = "https://mirrors.cloud.tencent.com/pypi/packages/df/b2/87e62e8c3e2f4b32e5fe99e0b86d576da1312593b39f47d8ceef365e95ed/packaging-26.2-py3-none-any.whl", hash = "sha256:5fc45236b9446107ff2415ce77c807cee2862cb6fac22b8a73826d0693b0980e" }, ] +[[package]] +name = "paho-mqtt" +version = "1.6.1" +source = { registry = "https://mirrors.cloud.tencent.com/pypi/simple/" } +sdist = { url = "https://mirrors.cloud.tencent.com/pypi/packages/f8/dd/4b75dcba025f8647bc9862ac17299e0d7d12d3beadbf026d8c8d74215c12/paho-mqtt-1.6.1.tar.gz", hash = "sha256:2a8291c81623aec00372b5a85558a372c747cbca8e9934dfe218638b8eefc26f" } + [[package]] name = "pineagents-demo-server" version = "0.1.0" @@ -809,6 +815,7 @@ dependencies = [ { name = "asyncmy" }, { name = "fastapi" }, { name = "httpx" }, + { name = "paho-mqtt" }, { name = "pydantic" }, { name = "pyjwt" }, { name = "python-multipart" }, @@ -831,6 +838,7 @@ requires-dist = [ { name = "asyncmy", specifier = ">=0.2.9" }, { name = "fastapi", specifier = ">=0.115.0" }, { name = "httpx", specifier = ">=0.27.0" }, + { name = "paho-mqtt", specifier = ">=1.6,<2" }, { name = "pydantic", specifier = ">=2.7.0" }, { name = "pyjwt", specifier = ">=2.8.0" }, { name = "python-multipart", specifier = ">=0.0.12" },