From 81f4dfd0ad3daba1f2dbe3942899e3f7dbfbf54b Mon Sep 17 00:00:00 2001 From: Pine Date: Mon, 24 Aug 2026 18:06:51 +0800 Subject: [PATCH] =?UTF-8?q?feat(park):=20=E5=9B=AD=E5=8C=BA=E5=A4=9A?= =?UTF-8?q?=E7=A7=9F=E6=88=B7=20=E2=80=94=20SQLite(tenants=E8=A1=A8+compan?= =?UTF-8?q?ies+kb)=20+=20=E6=AF=8F=E7=A7=9F=E6=88=B7=E5=BC=95=E6=93=8E=20+?= =?UTF-8?q?=20=E8=AE=A4=E8=AF=81=20+=20=E4=BD=9C=E7=94=A8=E5=9F=9F?= =?UTF-8?q?=E9=9A=94=E7=A6=BB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - tenants.py:park_tenants/park_companies/park_kb_docs 三表(serverdata/park/park.db), 每园区独立 name/intro(两段式)/auth(账密 salted hash)/data(大屏锚点)/agent/kb; ensure_default_tenant 自愈建默认园区+灌 39 家企业。 - sim_engine 重构为 SimEngine(data)——每租户独立实例,get_engine(tid)/refresh_engine(tid)。 - auth.py:/park/auth/login(账密→长效 tenant JWT)、/park/auth/me、require_tenant(401)。 - mqtt:register/publish_command/publish_tick 按 tenant_id 作用域(拓扑为 command/tid/clientId)。 - routers:/park/tenants CRUD + 全部数据/指令端点租户作用域(authorization token 或 ?tenant_id)。 - app.py lifespan init_db + 逐租户 tick/publish;tools._query_companies 读租户库。 - .gitignore 增 serverdata/park/。TestClient:建租户/登录/快照隔离/企业CRUD/指令鉴权(无token401) 全过。 --- .gitignore | 2 + app/park/app.py | 17 +- app/park/auth.py | 43 +++++ app/park/mqtt.py | 50 +++-- app/park/routers.py | 192 ++++++++++++++----- app/park/sim_engine.py | 293 ++++++++++++---------------- app/park/tenants.py | 419 +++++++++++++++++++++++++++++++++++++++++ app/park/tools.py | 6 +- 8 files changed, 779 insertions(+), 243 deletions(-) create mode 100644 app/park/auth.py create mode 100644 app/park/tenants.py diff --git a/.gitignore b/.gitignore index b7aefbb..dab7485 100644 --- a/.gitignore +++ b/.gitignore @@ -28,6 +28,8 @@ serverdata/keys/ serverdata/uploads/ serverdata/compute-data/ serverdata/compute-logs/ +# 园区子应用运行时数据(SQLite 多租户库/媒体/大屏数据,勿入库) +serverdata/park/ serverdata/mysql-data/ serverdata/redis-data/ serverdata/emqx-data/ diff --git a/app/park/app.py b/app/park/app.py index 4b3c92d..4eb4853 100644 --- a/app/park/app.py +++ b/app/park/app.py @@ -19,7 +19,8 @@ from .config import settings from .event_bus import bus from .mqtt import hub from .routers import router -from .sim_engine import sim_engine +from .sim_engine import get_engine, sim_engine +from . import tenants log = logging.getLogger("dpm.park") @@ -28,6 +29,11 @@ settings.MEDIA_DIR.mkdir(parents=True, exist_ok=True) @asynccontextmanager async def lifespan(app: FastAPI): + # SQLite 多租户数据初始化(建表 + 默认园区) + try: + tenants.init_db() + except Exception as e: # noqa: BLE001 + log.warning("园区 DB 初始化失败: %s", e) bus.set_loop(asyncio.get_running_loop()) hub.start() @@ -45,8 +51,13 @@ async def lifespan(app: FastAPI): def tick_loop(): while not stop.is_set(): try: - snap = sim_engine.tick() - hub.publish_tick(snap) + # 多租户:逐园区 tick + 按租户发布 + for t in tenants.list_tenants(): + try: + snap = get_engine(t["id"]).tick() + hub.publish_tick(snap, t["id"]) + except Exception as e: # noqa: BLE001 + log.warning("园区 %s tick 异常: %s", t.get("id"), e) except Exception as e: # noqa: BLE001 log.warning("园区数据 tick 异常: %s", e) stop.wait(settings.MQTT_TICK_INTERVAL) diff --git a/app/park/auth.py b/app/park/auth.py new file mode 100644 index 0000000..b527ed9 --- /dev/null +++ b/app/park/auth.py @@ -0,0 +1,43 @@ +# -*- coding: utf-8 -*- +"""园区租户认证 —— 大屏登录(账号密码 → 长效 tenant token,一次登录持久保持)。""" + +from __future__ import annotations + +import datetime +import jwt +from fastapi import Header, HTTPException + +from .config import settings + +_SECRET = settings.ADMIN_SECRET +_ALG = "HS256" +_TTL_DAYS = 365 + + +def create_token(tenant_id: str) -> str: + now = datetime.datetime.now(datetime.timezone.utc) + payload = { + "typ": "park", + "tid": tenant_id, + "iat": now, + "exp": now + datetime.timedelta(days=_TTL_DAYS), + } + return jwt.encode(payload, _SECRET, algorithm=_ALG) + + +def parse_token(token: str) -> str: + """返回 tenant_id;非法/过期抛 401。""" + try: + payload = jwt.decode(token, _SECRET, algorithms=[_ALG]) + except jwt.PyJWTError as e: # noqa: BLE001 + raise HTTPException(status_code=401, detail=f"park token 无效: {e}") from e + if payload.get("typ") != "park" or not payload.get("tid"): + raise HTTPException(status_code=401, detail="park token 无效") + return payload["tid"] + + +def require_tenant(authorization: str | None = Header(default=None, description="Bearer ")) -> str: + """大屏接口鉴权:从 Authorization 解析 tenant_id。""" + if not authorization or not authorization.startswith("Bearer "): + raise HTTPException(status_code=401, detail="缺少 park token") + return parse_token(authorization.split(" ", 1)[1]) diff --git a/app/park/mqtt.py b/app/park/mqtt.py index 81c5690..05ba879 100644 --- a/app/park/mqtt.py +++ b/app/park/mqtt.py @@ -87,13 +87,23 @@ class MqttHub: 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 = self.devices.get(cid, {"role": "", "parent_id": "", "tenant_id": ""}) cur["last_seen"] = time.time() + if payload.get("tenant_id"): + cur["tenant_id"] = payload["tenant_id"] self.devices[cid] = cur self._prune() except Exception: # noqa: BLE001 pass + def _online_ids(self, tenant_id: str = "") -> list[str]: + with self._lock: + cutoff = time.time() - settings.SCREEN_TTL + return [ + cid for cid, v in self.devices.items() + if v.get("last_seen", 0) > cutoff and (not tenant_id or v.get("tenant_id") == tenant_id) + ] + def _prune(self): """剔除超时未上报的设备""" cutoff = time.time() - settings.SCREEN_TTL @@ -135,8 +145,8 @@ class MqttHub: log.warning("MQTT 发布失败: %s", e) return False - def register_device(self, device_id, role="", parent_id=""): - """前端上报设备:唯一 device_id + 角色(main/secondary) + 从属主屏 parent_id。""" + def register_device(self, device_id, role="", parent_id="", tenant_id=""): + """前端上报设备:唯一 device_id + 角色(main/secondary) + 从属主屏 parent_id + 所属园区 tenant。""" device_id = (device_id or "").strip() if not device_id: return @@ -146,16 +156,18 @@ class MqttHub: cur["role"] = role if parent_id: cur["parent_id"] = parent_id + if tenant_id: + cur["tenant_id"] = tenant_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")) + log.info("mqtt: 设备注册 id=%s role=%s parent=%s tenant=%s", device_id, role or cur.get("role"), parent_id or cur.get("parent_id"), tenant_id) - def publish_command(self, action, params=None, screen_id=None, screen_role=None): - """按客户端(设备)路由下发指令: + def publish_command(self, action, params=None, screen_id=None, screen_role=None, tenant_id=""): + """按客户端(设备)路由下发指令(多租户:只遍历该租户在线设备): · screen_id 精确到某台设备 → 只发到该设备独有 topic · screen_role(main/secondary) → 解析到该角色的设备,发到其独有 topic - · 两者皆空(全部)→ 遍历在线设备列表,依次下发到各自独有 topic + · 两者皆空(全部)→ 遍历该租户在线设备,依次下发到各自独有 topic 客户端只订阅自己的独有频道,因此定向指令绝不会被其它屏幕收到。""" payload = { "cmd_id": new_cmd_id(), @@ -170,24 +182,22 @@ class MqttHub: elif screen_role in ("main", "secondary"): with self._lock: for cid, v in self.devices.items(): - if v.get("role") == screen_role: + if v.get("role") == screen_role and (not tenant_id or v.get("tenant_id") == tenant_id): target = cid break if target: # 定向:只发到目标设备独有频道 - return self._publish_to(payload, target, screen_role) + return self._publish_to(payload, target, screen_role, tenant_id) if requested: # 指定了目标但未找到对应设备:不广播,避免误发到全部 - log.warning("mqtt: 目标设备未找到(id=%s role=%s),未下发", screen_id, screen_role) + log.warning("mqtt: 目标设备未找到(id=%s role=%s tenant=%s),未下发", screen_id, screen_role, tenant_id) 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] + # 全部:遍历该租户在线设备,依次下发到各自独有频道 + online = self._online_ids(tenant_id) ok = True published_topics = [] for cid in online: @@ -195,7 +205,8 @@ class MqttHub: 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): + topic = f"{settings.TOPIC_COMMAND}/{tenant_id}/{cid}" if tenant_id else f"{settings.TOPIC_COMMAND}/{cid}" + if self.publish(topic, p): bus.emit(p) published_topics.append(cid) else: @@ -214,9 +225,9 @@ class MqttHub: log.info("mqtt: 全部下发 %s 到 %d 台设备", action, len(published_topics)) return payload - def _publish_to(self, payload, target, screen_role=""): + def _publish_to(self, payload, target, screen_role="", tenant_id=""): """发到指定设备独有频道,并记录 last_command。""" - topic = f"{settings.TOPIC_COMMAND}/{target}" + topic = f"{settings.TOPIC_COMMAND}/{tenant_id}/{target}" if tenant_id else f"{settings.TOPIC_COMMAND}/{target}" payload["screen_id"] = target with self._lock: payload["screen_role"] = self.devices.get(target, {}).get("role", screen_role or "") @@ -234,9 +245,10 @@ class MqttHub: bus.emit(payload) return payload - def publish_tick(self, snapshot): + def publish_tick(self, snapshot, tenant_id=""): payload = {"ts": int(time.time() * 1000), "snapshot": snapshot} - self.publish(settings.TOPIC_TICK, payload, qos=0) + topic = f"{settings.TOPIC_TICK}/{tenant_id}" if tenant_id else settings.TOPIC_TICK + self.publish(topic, payload, qos=0) hub = MqttHub() diff --git a/app/park/routers.py b/app/park/routers.py index 5e1cbb1..801ecb1 100644 --- a/app/park/routers.py +++ b/app/park/routers.py @@ -11,22 +11,100 @@ import json import logging import time -from fastapi import APIRouter, File, Request, UploadFile +from fastapi import APIRouter, File, Header, HTTPException, Request, UploadFile from fastapi.responses import JSONResponse, StreamingResponse from pydantic import BaseModel, Field -from . import park_config +from . import park_config, tenants +from .auth import create_token, parse_token, require_tenant from .config import settings from .event_bus import bus from .mqtt import hub -from .sim_engine import sim_engine +from .sim_engine import engine_data, get_engine, refresh_engine, sim_engine from .storage import storage +logger = logging.getLogger("dpm.api") + + +def _resolve_tenant(authorization: str | None, tenant_id: str | None) -> str: + """解析目标园区:大屏 park-token > ?tenant_id= > 默认园区。""" + if authorization and authorization.startswith("Bearer "): + return parse_token(authorization.split(" ", 1)[1]) + if tenant_id: + return tenant_id + return 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 + + +@router.post("/auth/login", summary="大屏/园区登录") +async def park_login(body: LoginBody): + t = tenants.verify_login(body.username, body.password) + if not t: + return JSONResponse({"ok": False, "error": "账号或密码不正确"}, status_code=401) + 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 = 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": tenants.list_tenants()} + + +@router.post("/tenants", summary="创建园区(名称/两段式简介/账密,seed 默认数据)") +async def tenant_create(body: TenantBody): + t = 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: TenantBody): + patch = {k: v for k, v in body.model_dump().items() if v is not None} + t = 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 = tenants.delete_tenant(tid) + return {"ok": ok, "id": tid} + + # ==================== 数据模型 ==================== class SettingsBody(BaseModel): @@ -135,97 +213,107 @@ async def runtime_config(request: Request): # ==================== 大屏数据(sim_engine) ==================== @router.get("/api/dashboard/snapshot") -async def dashboard_snapshot(): - return sim_engine.snapshot() +async def dashboard_snapshot(authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _resolve_tenant(authorization, tenant_id) + return get_engine(tid).snapshot() @router.get("/api/dashboard/overview") -async def dashboard_overview(): - return sim_engine.snapshot() +async def dashboard_overview(authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _resolve_tenant(authorization, tenant_id) + return get_engine(tid).snapshot() @router.get("/api/park/zones") -async def park_zones(): - snap = sim_engine.snapshot() - return snap.get("zones", []) +async def park_zones(authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _resolve_tenant(authorization, tenant_id) + return get_engine(tid).snapshot().get("zones", []) -# ==================== 入驻企业(park_config 主数据源) ==================== +# ==================== 入驻企业(租户 data.companies 主数据源) ==================== @router.get("/api/park/companies") -async def park_companies(): - return park_config.list_companies() +async def park_companies(authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _resolve_tenant(authorization, tenant_id) + return tenants.list_companies(tid) @router.post("/api/park/companies") -async def park_company_create(body: CompanyBody): - return park_config.create_company(body.model_dump()) +async def park_company_create(body: CompanyBody, authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _resolve_tenant(authorization, tenant_id) + return tenants.create_company(tid, body.model_dump()) @router.put("/api/park/companies/{cid}") -async def park_company_update(cid: str, body: CompanyBody): - row = park_config.update_company(cid, body.model_dump(exclude_unset=True)) +async def park_company_update(cid: str, body: CompanyBody, authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _resolve_tenant(authorization, tenant_id) + row = 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): - ok = park_config.delete_company(cid) +async def park_company_delete(cid: str, authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _resolve_tenant(authorization, tenant_id) + ok = tenants.delete_company(tid, cid) return {"ok": ok, "id": cid} # ==================== 园区智能体 / 知识库 / 大屏数据 ==================== @router.get("/api/agent/config") -async def agent_config_get(): - return park_config.get_agent() +async def agent_config_get(authorization: str | None = Header(None), tenant_id: str | None = None): + return tenants.get_agent(_resolve_tenant(authorization, tenant_id)) @router.put("/api/agent/config") -async def agent_config_put(body: AgentBody): - return park_config.update_agent(body.model_dump(exclude_none=True)) +async def agent_config_put(body: AgentBody, authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _resolve_tenant(authorization, tenant_id) + return tenants.update_agent(tid, body.model_dump(exclude_none=True)) @router.get("/api/kb/docs") -async def kb_docs_list(): - return park_config.list_kb_docs() +async def kb_docs_list(authorization: str | None = Header(None), tenant_id: str | None = None): + return tenants.kb_docs(_resolve_tenant(authorization, tenant_id)) @router.post("/api/kb/docs") -async def kb_docs_create(body: KbDocBody): - return park_config.create_kb_doc(body.model_dump()) +async def kb_docs_create(body: KbDocBody, authorization: str | None = Header(None), tenant_id: str | None = None): + return tenants.create_kb_doc(_resolve_tenant(authorization, tenant_id), body.model_dump()) @router.put("/api/kb/docs/{did}") -async def kb_docs_update(did: str, body: KbDocBody): - row = park_config.update_kb_doc(did, body.model_dump(exclude_unset=True)) +async def kb_docs_update(did: str, body: KbDocBody, authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _resolve_tenant(authorization, tenant_id) + row = 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): - ok = park_config.delete_kb_doc(did) +async def kb_docs_delete(did: str, authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _resolve_tenant(authorization, tenant_id) + ok = tenants.delete_kb_doc(tid, did) return {"ok": ok, "id": did} @router.post("/api/kb/reindex") -async def 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(): - return park_config.get_screen() +async def screen_data_get(authorization: str | None = Header(None), tenant_id: str | None = None): + return tenants.screen_view(_resolve_tenant(authorization, tenant_id)) @router.put("/api/screen/data") -async def screen_data_put(body: ScreenBody): - return park_config.update_screen(body.model_dump(exclude_none=True)) +async def screen_data_put(body: ScreenBody, authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _resolve_tenant(authorization, tenant_id) + return tenants.update_screen(tid, body.model_dump(exclude_none=True)) # ==================== AI / 知识库 / 语音识别 / 视觉(重模块已迁入) ==================== @@ -378,24 +466,40 @@ _VALID_DISPLAY_ACTIONS = { @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} +async def display_register(body: RegisterBody, authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _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): +async def display_command_publish(body: DisplayCommandBody, authorization: str | None = Header(None)): + tid = 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) + 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", - "mqtt_connected": hub.connected} + "tenant_id": tid, "mqtt_connected": hub.connected} @router.get("/api/display/state") -async def display_state(): - return hub.status() +async def display_state(authorization: str | None = Header(None), tenant_id: str | None = None): + tid = _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 不可用时前端回退) ==================== diff --git a/app/park/sim_engine.py b/app/park/sim_engine.py index d777048..a587fee 100644 --- a/app/park/sim_engine.py +++ b/app/park/sim_engine.py @@ -1,50 +1,15 @@ # -*- 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/工具仅做轻微演示波动 +"""园区数据引擎(后端)—— 真实数据 + 测算演示(每租户独立实例)。 + +数据锚点由租户 data 注入(可在管理端全量配置);引擎按 tenant_id 缓存独立实例, +实时 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 = [ "云南派音人工智能科技有限公司", "米勒克尔蓝宝石珠宝产业化(项目)", "中泰研学合作(项目)", "云南宸中低空经济有限公司", "昆明智海银高文化科技有限公司", "人工智能机器人大模型训练(项目)", @@ -62,31 +27,18 @@ COMPANIES = [ "蓝智科技(项目)", "云品出滇·纸享万家(项目)", "云迹轻旅—昆明本土一站式轻量化文旅服务工作室(项目)", "启元人工智能科技(项目)", ] - -# 宣传册四大孵化区域 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}, + {"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"}, @@ -99,50 +51,36 @@ FEED = [ {"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"}, ] +# ---------- 租户锚点默认数据(管理端可全量覆盖) ---------- +DEFAULT_DATA = { + "as_of": "2026-07", + "token": {"monthly_yi": 200, "daily_wan": 66666.7, "rate": 7700, "tools_per_call_tokens": 150000}, + "projects": {"inPark": 38, "capacity": 49, "invested": 560}, + "jobs": {"total": 213}, + "revenue": {"total": 252.15, "tax": 1.03, "today": 0}, + "cumulative": {"incubated": 181, "inIncubation": 31, "graduated": 149, "patents": 171, "jobs": 1529, "revenue": 21911.14, "tax": 426.2, "mentors": 50, "successRate": 95.91}, + "park": {"area": 3000, "founded": 2009, "provinceLevel": 2012, "address": "昆明市民航路229号(人力资源中心1-2楼)", "phone": "18087174891", "email": "714066050@qq.com"}, + "zones": ZONES, + "companies": [{"name": n, "zone": "", "room": "", "industry": "", "bio": "", "founder": ""} for n in COMPANIES], + "industryMix": INDUSTRY_MIX, + "feed": FEED, +} + 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 @@ -152,102 +90,70 @@ def _make_series(base, noise, n=30): 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, + "firstToken": round(1.3 + (r() - 0.5) * 0.2, 2), + "throughput": round(157 + (r() - 0.5) * 14), + "cacheHit": round(98 + (r() - 0.5) * 1.2, 1), + "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), + "steps": steps, } class SimEngine: - """带锁的快照引擎:单例供 API 与 MQTT tick 共用""" + """每租户独立快照引擎:按租户 data 锚点测算。""" - def __init__(self): + def __init__(self, data=None): + self.data = data or DEFAULT_DATA self._lock = threading.RLock() - self._snap = init_snapshot() + self._snap = self._init_snapshot() + + # ---- 测算 ---- + def _today_values(self, prev_noise=None): + token = self.data.get("token", {}) + rate = token.get("rate", 7700) + tpc = token.get("tools_per_call_tokens", 150000) + secs = _day_secs() + today_wan = round(secs * rate / 10000) + base_tools = secs * rate / 10000 * 10000 / tpc + noise = max(-10.0, min(10.0, (prev_noise or 0.0) + (random.random() - 0.5) * 6.0)) + return today_wan, max(0, round(base_tools + noise)), noise + + def _init_snapshot(self): + d = self.data + token = d.get("token", {}) + today_wan, today_tools, noise = self._today_values() + names = [c["name"] if isinstance(c, dict) else c for c in d.get("companies", [])] + return { + "as_of": d.get("as_of", "2026-07"), + "_toolsNoise": noise, + "live": _live_values(), + "token": { + "today": today_wan, + "total": token.get("monthly_yi", 200) * 10000 + today_wan, + "rate": token.get("rate", 7700), + "series": _make_series(token.get("daily_wan", 66666.7), token.get("daily_wan", 66666.7) * 0.16), + }, + "tools": { + "today": today_tools, + "total": round(token.get("monthly_yi", 200) * 100000000 / token.get("tools_per_call_tokens", 150000)) + today_tools, + "success": 98.5, + "nodes": 12, + }, + "projects": dict(d.get("projects", {})), + "jobs": dict(d.get("jobs", {})), + "revenue": dict(d.get("revenue", {})), + "cumulative": dict(d.get("cumulative", {})), + "park": dict(d.get("park", {})), + "zones": d.get("zones", ZONES), + "companies": names, + "industryMix": d.get("industryMix", INDUSTRY_MIX), + "feed": d.get("feed", FEED), + } def snapshot(self): with self._lock: @@ -255,8 +161,47 @@ class SimEngine: def tick(self): with self._lock: - self._snap = next_snapshot(self._snap) + d = self.data.get("token", {}) + today_wan, today_tools, noise = self._today_values(self._snap.get("_toolsNoise")) + self._snap.update({ + "_toolsNoise": noise, + "live": _live_values(self._snap.get("live")), + "token": { + "today": today_wan, + "total": d.get("monthly_yi", 200) * 10000 + today_wan, + "rate": d.get("rate", 7700), + "series": self._snap["token"]["series"][:-1] + [today_wan], + }, + "tools": { + "today": today_tools, + "total": round(d.get("monthly_yi", 200) * 100000000 / d.get("tools_per_call_tokens", 150000)) + today_tools, + "success": 98.5, + "nodes": 12, + }, + }) return dict(self._snap) -sim_engine = SimEngine() +# ---- 每租户引擎注册表 ---- +_engines: dict[str, SimEngine] = {} + +# 兼容:默认全局引擎(未带租户时回退) +sim_engine = SimEngine(DEFAULT_DATA) + + +def get_engine(tenant_id: str | None = None) -> SimEngine: + if not tenant_id: + return sim_engine + if tenant_id not in _engines: + from .tenants import get_tenant_data + _engines[tenant_id] = SimEngine(get_tenant_data(tenant_id)) + return _engines[tenant_id] + + +def refresh_engine(tenant_id: str) -> SimEngine: + _engines.pop(tenant_id, None) + return get_engine(tenant_id) + + +def engine_data(tenant_id: str | None = None) -> dict: + return get_engine(tenant_id).data diff --git a/app/park/tenants.py b/app/park/tenants.py new file mode 100644 index 0000000..cd9c407 --- /dev/null +++ b/app/park/tenants.py @@ -0,0 +1,419 @@ +# -*- coding: utf-8 -*- +"""园区租户(多租户)—— SQLite 持久化(serverdata/park/park.db,后续可迁 MySQL/等)。 + +表: + park_tenants(id, name, intro_json, username, salt, password_hash, data_json, agent_json, created_at) + park_companies(id, tenant_id, name, zone, room, industry, bio, founder, status, employees, created_at) + park_kb_docs(id, tenant_id, grp, title, content_md, created_at) + +data_json(大屏锚点:as_of/token/projects/jobs/revenue/cumulative/park/zones/industryMix/feed) +与 agent_json 用 JSON 列(灵活、易迁移);companies/kb 用表(可直接查询/CRUD)。 +""" +from __future__ import annotations + +import hashlib +import json +import secrets +import sqlite3 +import time +import uuid +from pathlib import Path + +from .config import PARK_DIR +from .sim_engine import DEFAULT_DATA, refresh_engine + +DB_PATH = Path(PARK_DIR) / "park.db" +_DEFAULT_TENANT_ID = "T001" + +SCHEMA = """ +CREATE TABLE IF NOT EXISTS park_tenants ( + id TEXT PRIMARY KEY, + name TEXT, + intro_json TEXT, + username TEXT, + salt TEXT, + password_hash TEXT, + data_json TEXT, + agent_json TEXT, + created_at TEXT +); +CREATE TABLE IF NOT EXISTS park_companies ( + id TEXT PRIMARY KEY, + tenant_id TEXT, + name TEXT, + zone TEXT DEFAULT '', + room TEXT DEFAULT '', + industry TEXT DEFAULT '', + bio TEXT DEFAULT '', + founder TEXT DEFAULT '', + status TEXT DEFAULT 'applying', + employees INTEGER, + created_at TEXT +); +CREATE TABLE IF NOT EXISTS park_kb_docs ( + id TEXT PRIMARY KEY, + tenant_id TEXT, + grp TEXT DEFAULT 'general', + title TEXT, + content_md TEXT DEFAULT '', + created_at TEXT +); +""" + + +def _conn() -> sqlite3.Connection: + DB_PATH.parent.mkdir(parents=True, exist_ok=True) + conn = sqlite3.connect(DB_PATH, check_same_thread=False) + conn.row_factory = sqlite3.Row + conn.execute("PRAGMA journal_mode=WAL") + return conn + + +def init_db(): + conn = _conn() + conn.executescript(SCHEMA) + conn.commit() + conn.close() + ensure_default_tenant() + + +def _hash_password(password: str, salt: str | None = None) -> tuple[str, str]: + salt = salt or secrets.token_hex(16) + h = hashlib.pbkdf2_hmac("sha256", password.encode(), salt.encode(), 100_000).hex() + return salt, h + + +def _verify_password(password: str, salt: str, expected: str) -> bool: + return _hash_password(password, salt)[1] == expected + + +def _now() -> str: + return time.strftime("%Y-%m-%d %H:%M:%S") + + +def _new_id(prefix: str) -> str: + return f"{prefix}{int(uuid.uuid4().int % 1000000000):09d}" + + +# ---------------- 租户 ---------------- + +def _row_to_tenant(row) -> dict: + return { + "id": row["id"], + "name": row["name"], + "intro": json.loads(row["intro_json"] or "[]"), + "auth": {"username": row["username"], "salt": row["salt"], "password_hash": row["password_hash"]}, + "data": json.loads(row["data_json"] or "{}"), + "agent": json.loads(row["agent_json"] or "{}"), + } + + +def list_tenants() -> list[dict]: + conn = _conn() + rows = conn.execute("SELECT * FROM park_tenants ORDER BY created_at DESC").fetchall() + conn.close() + return [_row_to_tenant(r) for r in rows] + + +def get_tenant(tenant_id: str) -> dict | None: + conn = _conn() + row = conn.execute("SELECT * FROM park_tenants WHERE id=?", (tenant_id,)).fetchone() + conn.close() + return _row_to_tenant(row) if row else None + + +def _get_data(tenant_id: str) -> dict: + t = get_tenant(tenant_id) + return (t or {}).get("data", {}) + + +def get_tenant_data(tenant_id: str) -> dict: + """给 sim_engine 用:data 锚点 + 注入 companies(来自 park_companies 表)。""" + data = _get_data(tenant_id) + comps = list_companies(tenant_id) + if comps: + data = dict(data) + data["companies"] = comps + return data or dict(DEFAULT_DATA) + + +def create_tenant(name: str, intro: list[str], username: str, password: str) -> dict: + salt, h = _hash_password(password) + tid = _new_id("T") + data = dict(DEFAULT_DATA) + conn = _conn() + conn.execute( + "INSERT INTO park_tenants (id,name,intro_json,username,salt,password_hash,data_json,agent_json,created_at) VALUES (?,?,?,?,?,?,?,?,?)", + (tid, name, json.dumps(intro or ["", ""]), username, salt, h, json.dumps(data, ensure_ascii=False), json.dumps({}, ensure_ascii=False), _now()), + ) + conn.commit() + conn.close() + _seed_companies(tid, DEFAULT_DATA["companies"]) + refresh_engine(tid) + return _summary(tid, name, intro, username) + + +def update_tenant(tenant_id: str, patch: dict) -> dict | None: + conn = _conn() + t = get_tenant(tenant_id) + if t is None: + conn.close() + return None + sets, args = [], [] + if patch.get("name") is not None: + sets.append("name=?"); args.append(patch["name"]) + if patch.get("intro") is not None: + sets.append("intro_json=?"); args.append(json.dumps(patch["intro"], ensure_ascii=False)) + if patch.get("data") is not None: + sets.append("data_json=?"); args.append(json.dumps(patch["data"], ensure_ascii=False)) + if patch.get("agent") is not None: + sets.append("agent_json=?"); args.append(json.dumps(patch["agent"], ensure_ascii=False)) + if patch.get("username") is not None: + sets.append("username=?"); args.append(patch["username"]) + if patch.get("password"): + salt, h = _hash_password(patch["password"]) + sets.append("salt=?"); sets.append("password_hash=?"); args += [salt, h] + if not sets: + conn.close() + return _summary(tenant_id, t["name"], t["intro"], t["auth"]["username"]) + args.append(tenant_id) + conn.execute(f"UPDATE park_tenants SET {', '.join(sets)} WHERE id=?", args) + conn.commit() + conn.close() + if patch.get("data"): + refresh_engine(tenant_id) + t2 = get_tenant(tenant_id) + return _summary(tenant_id, t2["name"], t2["intro"], t2["auth"]["username"]) + + +def delete_tenant(tenant_id: str) -> bool: + conn = _conn() + cur = conn.execute("DELETE FROM park_tenants WHERE id=?", (tenant_id,)) + conn.execute("DELETE FROM park_companies WHERE tenant_id=?", (tenant_id,)) + conn.execute("DELETE FROM park_kb_docs WHERE tenant_id=?", (tenant_id,)) + conn.commit() + conn.close() + refresh_engine(tenant_id) # 淘汰缓存引擎 + from .sim_engine import _engines + _engines.pop(tenant_id, None) + return cur.rowcount > 0 + + +def verify_login(username: str, password: str) -> dict | None: + conn = _conn() + row = conn.execute("SELECT * FROM park_tenants WHERE username=?", (username,)).fetchone() + conn.close() + if row and _verify_password(password, row["salt"], row["password_hash"]): + return _summary(row["id"], row["name"], json.loads(row["intro_json"] or "[]"), row["username"]) + return None + + +def ensure_default_tenant() -> str: + """确保至少存在一个默认园区(昆明市大学生创业园)。已有则补种企业(自愈)。""" + conn = _conn() + if conn.execute("SELECT 1 FROM park_tenants WHERE id=?", (_DEFAULT_TENANT_ID,)).fetchone(): + conn.close() + if not list_companies(_DEFAULT_TENANT_ID): + _seed_companies(_DEFAULT_TENANT_ID, DEFAULT_DATA["companies"]) + return _DEFAULT_TENANT_ID + conn.close() + salt, h = _hash_password("123456") + data = dict(DEFAULT_DATA) + conn = _conn() + conn.execute( + "INSERT INTO park_tenants (id,name,intro_json,username,salt,password_hash,data_json,agent_json,created_at) VALUES (?,?,?,?,?,?,?,?,?)", + (_DEFAULT_TENANT_ID, "昆明市大学生创业园", json.dumps(["云南省首家政府主办大学生创业孵化园区", "空间 + 孵化 + 融资 + 政策 + AI 赋能 + 综合服务"], ensure_ascii=False), "admin", salt, h, json.dumps(data, ensure_ascii=False), json.dumps({}, ensure_ascii=False), _now()), + ) + conn.commit() + conn.close() + _seed_companies(_DEFAULT_TENANT_ID, DEFAULT_DATA["companies"]) + refresh_engine(_DEFAULT_TENANT_ID) + return _DEFAULT_TENANT_ID + + +def _summary(tid: str, name: str, intro: list, username: str) -> dict: + return {"id": tid, "name": name, "intro": intro, "username": username} + + +def _seed_companies(tenant_id: str, companies: list) -> None: + """把默认企业灌入 park_companies 表(新园区初始名录),管理端在其上改。""" + conn = _conn() + for c in companies: + name = c["name"] if isinstance(c, dict) else c + conn.execute( + "INSERT OR IGNORE INTO park_companies (id,tenant_id,name,zone,room,industry,bio,founder,status,employees,created_at) VALUES (?,?,?,?,?,?,?,?,?,?,?)", + (_new_id("PC"), tenant_id, name, c.get("zone", "") if isinstance(c, dict) else "", c.get("room", "") if isinstance(c, dict) else "", c.get("industry", "") if isinstance(c, dict) else "", c.get("bio", "") if isinstance(c, dict) else "", c.get("founder", "") if isinstance(c, dict) else "", "active", None, _now()), + ) + conn.commit() + conn.close() + + +def iter_tenants(): + for t in list_tenants(): + yield t["id"], t + + +# ---------------- 入驻企业(表级 CRUD) ---------------- + +def list_companies(tenant_id: str) -> list[dict]: + conn = _conn() + rows = conn.execute("SELECT * FROM park_companies WHERE tenant_id=? ORDER BY created_at DESC", (tenant_id,)).fetchall() + conn.close() + return [dict(r) for r in rows] + + +def create_company(tenant_id: str, payload: dict) -> dict: + row = { + "id": _new_id("PC"), "tenant_id": tenant_id, "name": payload.get("name", ""), + "zone": payload.get("zone", ""), "room": payload.get("room", ""), "industry": payload.get("industry", ""), + "bio": payload.get("bio", ""), "founder": payload.get("founder", ""), + "status": payload.get("status", "applying"), "employees": payload.get("employees"), "created_at": _now(), + } + conn = _conn() + conn.execute( + "INSERT INTO park_companies (id,tenant_id,name,zone,room,industry,bio,founder,status,employees,created_at) VALUES (?,?,?,?,?,?,?,?,?,?,?)", + (row["id"], row["tenant_id"], row["name"], row["zone"], row["room"], row["industry"], row["bio"], row["founder"], row["status"], row["employees"], row["created_at"]), + ) + conn.commit() + conn.close() + refresh_engine(tenant_id) + return {k: v for k, v in row.items() if k != "tenant_id"} + + +def update_company(tenant_id: str, cid: str, payload: dict) -> dict | None: + fields = {k: v for k, v in payload.items() if v is not None and k in ("name", "zone", "room", "industry", "bio", "founder", "status", "employees")} + if not fields: + return None + sets = ", ".join(f"{k}=?" for k in fields) + conn = _conn() + cur = conn.execute(f"UPDATE park_companies SET {sets} WHERE id=? AND tenant_id=?", list(fields.values()) + [cid, tenant_id]) + conn.commit() + conn.close() + refresh_engine(tenant_id) + if cur.rowcount == 0: + return None + conn = _conn() + row = conn.execute("SELECT * FROM park_companies WHERE id=? AND tenant_id=?", (cid, tenant_id)).fetchone() + conn.close() + return {k: v for k, v in dict(row).items() if k != "tenant_id"} + + +def delete_company(tenant_id: str, cid: str) -> bool: + conn = _conn() + cur = conn.execute("DELETE FROM park_companies WHERE id=? AND tenant_id=?", (cid, tenant_id)) + conn.commit() + conn.close() + refresh_engine(tenant_id) + return cur.rowcount > 0 + + +# ---------------- 园区智能体 ---------------- + +def get_agent(tenant_id: str) -> dict: + return get_tenant(tenant_id).get("agent", {}) + + +def update_agent(tenant_id: str, payload: dict) -> dict: + t = get_tenant(tenant_id) + agent = t.get("agent", {}) + for k, v in payload.items(): + if v is not None: + agent[k] = v + conn = _conn() + conn.execute("UPDATE park_tenants SET agent_json=? WHERE id=?", (json.dumps(agent, ensure_ascii=False), tenant_id)) + conn.commit() + conn.close() + return agent + + +# ---------------- 园区知识库(表级 CRUD) ---------------- + +def kb_docs(tenant_id: str) -> list[dict]: + conn = _conn() + rows = conn.execute("SELECT * FROM park_kb_docs WHERE tenant_id=? ORDER BY created_at DESC", (tenant_id,)).fetchall() + conn.close() + return [dict(r) for r in rows] + + +def create_kb_doc(tenant_id: str, payload: dict) -> dict: + row = {"id": _new_id("KB"), "tenant_id": tenant_id, "grp": payload.get("group", "general"), "title": payload.get("title", ""), "content_md": payload.get("content_md", ""), "created_at": _now()} + conn = _conn() + conn.execute("INSERT INTO park_kb_docs (id,tenant_id,grp,title,content_md,created_at) VALUES (?,?,?,?,?,?)", (row["id"], row["tenant_id"], row["grp"], row["title"], row["content_md"], row["created_at"])) + conn.commit() + conn.close() + return {k: v for k, v in row.items() if k != "tenant_id"} + + +def update_kb_doc(tenant_id: str, did: str, payload: dict) -> dict | None: + fields = {k: v for k, v in payload.items() if v is not None and k in ("grp", "title", "content_md")} + if not fields: + return get_kb_doc(tenant_id, did) + sets = ", ".join(f"{k}=?" for k in fields) + conn = _conn() + conn.execute(f"UPDATE park_kb_docs SET {sets} WHERE id=? AND tenant_id=?", list(fields.values()) + [did, tenant_id]) + conn.commit() + conn.close() + return get_kb_doc(tenant_id, did) + + +def delete_kb_doc(tenant_id: str, did: str) -> bool: + conn = _conn() + cur = conn.execute("DELETE FROM park_kb_docs WHERE id=? AND tenant_id=?", (did, tenant_id)) + conn.commit() + conn.close() + return cur.rowcount > 0 + + +def get_kb_doc(tenant_id: str, did: str) -> dict | None: + conn = _conn() + row = conn.execute("SELECT * FROM park_kb_docs WHERE id=? AND tenant_id=?", (did, tenant_id)).fetchone() + conn.close() + return {k: v for k, v in dict(row).items() if k != "tenant_id"} if row else None + + +# ---------------- 大屏数据(名称/简介 + 锚点) ---------------- + +def screen_view(tenant_id: str) -> dict: + t = get_tenant(tenant_id) or {} + d = t.get("data", {}) + park = d.get("park", {}) or {} + return { + "tenant_id": tenant_id, + "name": t.get("name", ""), + "intro": t.get("intro", []), + "founded": park.get("founded"), "province_level": park.get("province_level") or park.get("provinceLevel"), "area": park.get("area"), + "address": park.get("address"), "phone": park.get("phone"), "email": park.get("email"), + "capacity": d.get("projects", {}).get("capacity"), "invested": d.get("projects", {}).get("invested"), + "jobs": d.get("jobs", {}).get("total"), "revenue_total": d.get("revenue", {}).get("total"), "revenue_tax": d.get("revenue", {}).get("tax"), + } + + +def update_screen(tenant_id: str, payload: dict) -> dict: + t = get_tenant(tenant_id) or {} + d = dict(t.get("data", {})) + if payload.get("name"): + d["_name"] = payload["name"] # 名称存租户 name 列,见下 + park = dict(d.get("park", {})) + for k in ("founded", "area", "address", "phone", "email"): + if payload.get(k) is not None: + park[k] = payload[k] + d["park"] = park + for k, sub in (("capacity", "projects"), ("invested", "projects"), ("jobs", "jobs"), ("revenue_total", "revenue"), ("revenue_tax", "revenue")): + if payload.get(k) is not None: + s = dict(d.get(sub, {})) + target = "total" if sub in ("jobs", "revenue") and k != "revenue_total" and k != "revenue_tax" else ("capacity" if k == "capacity" else "invested" if k == "invested" else ("tax" if k == "revenue_tax" else "total")) + if sub == "revenue": + target = "total" if k == "revenue_total" else "tax" + if sub == "jobs": + target = "total" + s[target] = payload[k] + d[sub] = s + conn = _conn() + name = payload.get("name") or t["name"] + intro = payload.get("intro") or t.get("intro", []) + conn.execute("UPDATE park_tenants SET name=?, intro_json=?, data_json=? WHERE id=?", + (name, json.dumps(intro, ensure_ascii=False), json.dumps(d, ensure_ascii=False), tenant_id)) + conn.commit() + conn.close() + refresh_engine(tenant_id) + return screen_view(tenant_id) diff --git a/app/park/tools.py b/app/park/tools.py index e011e2d..7558ddb 100644 --- a/app/park/tools.py +++ b/app/park/tools.py @@ -49,10 +49,10 @@ def _get_park_overview(args): def _query_companies(args): - """园区入驻企业名录查询(关键词过滤,取自 park_config 主数据源)""" - from . import park_config + """园区入驻企业名录查询(关键词过滤,取自租户 SQLite 主数据源)""" + from . import tenants q = (args.get("keyword") or "").strip() - all_comps = park_config.list_companies() + all_comps = tenants.list_companies(tenants.ensure_default_tenant()) names = [c["name"] for c in all_comps] names = [n for n in names if q in n] if q else names return json.dumps({