From 416ade48427269458c4ee28cde0e18744b2f0bd8 Mon Sep 17 00:00:00 2001 From: Pine Date: Fri, 11 Sep 2026 21:15:36 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E7=BB=9F=E4=B8=80=E6=9D=83=E9=99=90?= =?UTF-8?q?=E8=81=9A=E5=90=88=E7=AB=AF=E7=82=B9=20+=20=E4=BA=8B=E4=BB=B6?= =?UTF-8?q?=E5=90=8D=E5=90=8C=E6=AD=A5=E6=98=B5=E7=A7=B0=20+=20IM/?= =?UTF-8?q?=E4=BC=81=E4=B8=9A=E8=B7=AF=E7=94=B1=E8=A1=A5=E5=BC=BA?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 新增 rbac_permissions 统一权限/组织归属聚合端点 - sync_event_name_to_nickname 脚本:事件名同步用户昵称 - im router/client、rbac_enterprise/opc/org/public 增强 - compute_catalog、nginx 配置、env.example 更新 --- app/api/routers/rbac_desktop_updates.py | 275 ++++++++++++++++-- app/api/routers/rbac_enterprise.py | 27 ++ app/api/routers/rbac_opc.py | 20 +- app/api/routers/rbac_org.py | 41 ++- app/api/routers/rbac_permissions.py | 78 +++++ app/api/routers/rbac_public.py | 61 ++++ app/im/client.py | 13 + app/im/router.py | 12 + app/main.py | 2 + app/services/compute_catalog.py | 15 + nginx/nginx/conf.d/opc.pinesound.cn.conf | 2 +- scripts/.sync_event_name_state.json | 5 + scripts/import_kunming_park.py | 2 +- scripts/recover_from_binlog.py | 2 +- scripts/sync_event_name_to_nickname.py | 232 +++++++++++++++ .../backup/opc_full_20260908_175552.sql | 2 +- serverrun/.env.example | 12 +- 17 files changed, 760 insertions(+), 41 deletions(-) create mode 100644 app/api/routers/rbac_permissions.py create mode 100644 scripts/.sync_event_name_state.json create mode 100644 scripts/sync_event_name_to_nickname.py diff --git a/app/api/routers/rbac_desktop_updates.py b/app/api/routers/rbac_desktop_updates.py index 5cb5939..e9f5770 100644 --- a/app/api/routers/rbac_desktop_updates.py +++ b/app/api/routers/rbac_desktop_updates.py @@ -21,7 +21,7 @@ import re import secrets from pathlib import Path -from fastapi import APIRouter, Depends, File, Form, HTTPException, Request, UploadFile +from fastapi import APIRouter, Depends, File, Form, HTTPException, Query, Request, UploadFile from fastapi.responses import FileResponse, RedirectResponse from pydantic import BaseModel, Field @@ -41,6 +41,43 @@ DESKTOP_TARGETS = ( "windows-x86_64", "linux-x86_64", ) + +# Tauri v2 可能携带的 target 格式映射(带 bundle 类型后缀) +# 例如 darwin-aarch64-app -> darwin-aarch64, windows-x86_64-nsis -> windows-x86_64 +_TARGET_ALIASES = { + "darwin-aarch64-app": "darwin-aarch64", + "darwin-aarch64-dmg": "darwin-aarch64", + "darwin-x86_64-app": "darwin-x86_64", + "darwin-x86_64-dmg": "darwin-x86_64", + "windows-x86_64-nsis": "windows-x86_64", + "windows-x86_64-msi": "windows-x86_64", + "linux-x86_64-appimage": "linux-x86_64", + "linux-x86_64-deb": "linux-x86_64", + "linux-x86_64-rpm": "linux-x86_64", +} + + +def _normalize_target(target: str) -> str: + """将 Tauri v2 可能携带的各种 target 格式规范化为标准格式。 + + 例如: + darwin-aarch64-app -> darwin-aarch64 + windows-x86_64-nsis -> windows-x86_64 + """ + target = (target or "").strip().lower() + if not target: + return "" + # 直接匹配 + if target in DESKTOP_TARGETS: + return target + # 别名映射 + if target in _TARGET_ALIASES: + return _TARGET_ALIASES[target] + # 尝试去掉最后一个后缀(如 -app, -nsis) + parts = target.rsplit("-", 1) + if len(parts) == 2 and parts[0] in DESKTOP_TARGETS: + return parts[0] + return target VERSION_RE = re.compile(r"^\d+\.\d+\.\d+([-+][0-9A-Za-z.\-]+)?$") MAX_ARTIFACT_BYTES = 2 * 1024 * 1024 * 1024 # 2GB(桌面端含 PyInstaller 后端,体积较大) @@ -246,18 +283,43 @@ def _client_version(request: Request) -> str | None: return value or None -def _build_manifest(version: dict, base_url: str) -> dict: - """组装 Tauri updater 标准 manifest(version/notes/pub_date/platforms)。""" +def _build_manifest(version: dict, base_url: str, target: str | None = None) -> dict: + """组装 Tauri updater 标准 manifest(version/notes/pub_date/platforms)。 + + Args: + version: 版本记录 + base_url: 基础 URL + target: 如果指定,只包含该平台的 artifacts;否则包含所有平台 + """ platforms: dict[str, dict] = {} - for target, art in (version.get("artifacts") or {}).items(): - url = (art or {}).get("url", "") - signature = (art or {}).get("signature", "") - if not url or not signature: - continue - platforms[target] = { - "url": f"{base_url}{url}", - "signature": signature, - } + artifacts = version.get("artifacts") or {} + + if target: + # 指定了 target,使用 _has_target_artifact 找到实际匹配的 target key + # (支持 darwin-aarch64-app 等带后缀的格式) + has, actual_key = _has_target_artifact(version, target) + if has and actual_key: + art = artifacts.get(actual_key) or {} + url = art.get("url", "") + signature = art.get("signature", "") + if url: + # manifest 中使用规范化后的 target(darwin-aarch64),Tauri 能识别 + platforms[target] = { + "url": f"{base_url}{url}", + "signature": signature or "", + } + else: + # 未指定 target,包含所有平台 + for t, art in artifacts.items(): + url = (art or {}).get("url", "") + signature = (art or {}).get("signature", "") + if not url: + continue + platforms[t] = { + "url": f"{base_url}{url}", + "signature": signature or "", + } + return { "version": version["version"], "notes": version.get("notes", "") or "", @@ -266,6 +328,108 @@ def _build_manifest(version: dict, base_url: str) -> dict: } +def _has_target_artifact(version: dict, target: str) -> tuple[bool, str | None]: + """检查版本是否包含指定平台的 artifacts,返回 (是否包含, 实际匹配的 target key)。 + + 支持多种 target 格式匹配: + - 精确匹配:darwin-aarch64 + - 带后缀匹配:darwin-aarch64-app, darwin-aarch64-dmg + - 反向匹配:如果 artifacts 中的 key 是 darwin-aarch64-app,target 是 darwin-aarch64,也能匹配 + """ + artifacts = version.get("artifacts") or {} + if not artifacts: + return False, None + + # 1. 精确匹配 + if target in artifacts: + art = artifacts[target] or {} + if art.get("url"): + return True, target + + # 2. 尝试带后缀的格式(target + 后缀) + for suffix in ("-app", "-dmg", "-nsis", "-msi", "-appimage", "-deb", "-rpm"): + key = f"{target}{suffix}" + if key in artifacts: + art = artifacts[key] or {} + if art.get("url"): + return True, key + + # 3. 反向匹配:artifacts 中的 key 可能带后缀,去掉后缀后匹配 + for key in artifacts: + if not key: + continue + # 去掉最后一个后缀(如 -app, -nsis) + parts = key.rsplit("-", 1) + if len(parts) == 2 and parts[0] == target: + art = artifacts[key] or {} + if art.get("url"): + return True, key + + return False, None + + +def _select_latest_for_target( + candidates: list[dict], + target: str, + client_version: str | None, + channel: str = "stable", +) -> dict | None: + """为指定平台选择最新可用版本。 + + 遍历所有候选版本,找到包含该平台 artifacts 的最新版本。 + 支持灰度逻辑:灰度未命中时退回上一全量版本。 + + 注意:先尝试筛选指定 channel 的版本,如果找不到该平台的 artifacts, + 就回退到所有版本(不筛选 channel)。这样即使最新版本的 channel + 不同,只要包含该平台的 artifacts,就能被找到。 + + Args: + candidates: 所有可发布版本(按版本从高到低排序) + target: 平台 target(如 darwin-aarch64) + client_version: 客户端版本号(用于灰度判定) + channel: 发布渠道(stable/beta) + + Returns: + 该平台的最新版本记录,或 None + """ + def _find_in_versions(versions: list[dict]) -> list[dict]: + """在指定版本列表中找到包含该平台 artifacts 的所有版本。""" + result = [] + for v in versions: + has, _ = _has_target_artifact(v, target) + if has: + result.append(v) + return result + + # 1. 先尝试筛选指定 channel 的版本 + channel_versions = [v for v in candidates if (v.get("channel") or "stable") == channel] + target_versions = _find_in_versions(channel_versions) + + # 2. 如果在指定 channel 中找不到,回退到所有版本(不筛选 channel) + if not target_versions: + target_versions = _find_in_versions(candidates) + + if not target_versions: + return None + + latest = target_versions[0] + + # 灰度逻辑:如果最新版本是灰度且未命中,退回上一全量版本 + if ( + latest.get("status") == "rolling" + and int(latest.get("rollout") or 0) < 100 + and client_version + ): + hit = int(hashlib.sha256(client_version.encode()).hexdigest(), 16) % 100 + if hit >= int(latest.get("rollout") or 0): + for v in target_versions[1:]: + if v.get("status") == "published": + return v + return None + + return latest + + def _select_latest( candidates: list[dict], client_version: str | None, @@ -780,20 +944,97 @@ async def delete_desktop_version( # 公开端点(桌面端更新源) # --------------------------------------------------------------------------- -@public_router.get("/desktop-updates/latest.json", summary="Tauri updater 更新清单(支持灰度)") -async def desktop_updates_latest(request: Request, db: Database = Depends(get_db)): +@public_router.get("/desktop-updates/latest.json", summary="Tauri updater 更新清单(支持灰度 + 按平台选择最新)") +async def desktop_updates_latest( + request: Request, + target: str = Query("", description="Tauri updater 自动携带的平台 target,如 darwin-aarch64"), + db: Database = Depends(get_db), +): + """Tauri updater 更新清单。 + + 支持两种模式: + 1. 指定 target(Tauri v2 自动携带):返回该平台的最新可用版本 + - 例如 macOS 用户请求 target=darwin-aarch64,返回 macOS 平台的最新版本 + - 即使 Windows 有更新的版本,macOS 用户也只看到 macOS 的最新版本 + 2. 未指定 target(兼容旧版本):返回全局最新版本的所有平台 artifacts + + 版本号比较由 Tauri 客户端完成,服务端只负责返回对应平台的最新可用版本。 + """ client_version = _client_version(request) channel = request.headers.get("X-App-Channel", "stable").strip() or "stable" candidates = await db.desktop_versions.publishable() - latest = _select_latest(candidates, client_version, channel) - if latest is None: + + # 如果指定了 target,按平台选择最新版本 + # Tauri v2 可能携带 darwin-aarch64-app 等带后缀的格式,需要规范化 + target = _normalize_target(target) + if target and target in DESKTOP_TARGETS: + latest = _select_latest_for_target(candidates, target, client_version, channel) + if latest is None: + # 该平台没有可用更新,但 Tauri updater 要求 platforms 中必须包含当前平台的 key, + # 否则会报错 "None of the fallback platforms were found in the response platforms object"。 + # 返回版本号 0.0.0(小于任何真实版本),Tauri 会判定为无更新。 + return { + "version": "0.0.0", + "notes": "", + "pub_date": utcnow_iso(), + "platforms": { + target: { + "url": "", + "signature": "", + } + }, + } + # 只返回该平台的 artifacts,version 字段填该平台的最新版本号 + return _build_manifest(latest, _base_url(request), target=target) + + # 未指定 target:返回每个平台各自的最新版本 + # 这样即使某个平台(如 macOS)没有全局最新版本的包,也能拿到该平台的最新可用版本 + # 例如:Windows 最新是 b3,macOS 最新是 b2,则返回的 platforms 中 + # windows-x86_64 对应 b3 的包,darwin-aarch64 对应 b2 的包 + base_url = _base_url(request) + platforms: dict[str, dict] = {} + version_records: list[dict] = [] + + for t in DESKTOP_TARGETS: + latest_for_target = _select_latest_for_target(candidates, t, client_version, channel) + if latest_for_target is None: + continue + # 获取该平台的 artifacts(支持带后缀的格式匹配) + has, actual_key = _has_target_artifact(latest_for_target, t) + if has and actual_key: + art = (latest_for_target.get("artifacts") or {}).get(actual_key) or {} + url = art.get("url", "") + signature = art.get("signature", "") + if url: + platforms[t] = { + "url": f"{base_url}{url}", + "signature": signature or "", + } + version_records.append(latest_for_target) + + if not platforms: return { "version": "0.0.0", "notes": "", "pub_date": utcnow_iso(), "platforms": {}, } - return _build_manifest(latest, _base_url(request)) + + # version 字段填所有平台中最高的版本号 + # Tauri updater 会比较这个版本号和当前版本号,如果更高则下载对应平台的包 + try: + from packaging.version import Version + highest_record = max(version_records, key=lambda v: Version(v["version"])) + except Exception: + # 回退:按字符串排序 + highest_record = max(version_records, key=lambda v: v["version"]) + + return { + "version": highest_record["version"], + "notes": highest_record.get("notes", "") or "", + "pub_date": highest_record.get("pub_date") or utcnow_iso(), + "platforms": platforms, + } @public_router.get("/desktop-updates/meta.json", summary="更新元信息(强制更新判定)") diff --git a/app/api/routers/rbac_enterprise.py b/app/api/routers/rbac_enterprise.py index 0298d76..ad4c68c 100644 --- a/app/api/routers/rbac_enterprise.py +++ b/app/api/routers/rbac_enterprise.py @@ -9,6 +9,8 @@ """ from __future__ import annotations +import logging + from fastapi import APIRouter, Depends, HTTPException, Request from pydantic import BaseModel from sqlalchemy import select @@ -19,9 +21,30 @@ from ...infrastructure.models import ParkCompany from ...infrastructure.repositories import Database from app.park import tenants as tnt +log = logging.getLogger("rbac.enterprise") + router = APIRouter(prefix="/admin/ent", tags=["enterprise"]) +async def _enterprise_sync_task(cid: str) -> None: + """查企业名后同步企业群;任何异常都吞掉,不影响主流程。""" + try: + from ...im import client as im_client + co = await tnt.get_company(cid) + await im_client.sync_enterprise_group(cid, (co or {}).get("name", "")) + except Exception: # noqa: BLE001 + log.warning("企业群同步失败 company=%s", cid, exc_info=True) + + +def _fire_enterprise_sync(cid: str) -> None: + """fire-and-forget 触发企业群同步(不阻断主流程,不等待结果)。""" + import asyncio + try: + asyncio.create_task(_enterprise_sync_task(cid)) + except Exception: # noqa: BLE001 + pass + + async def _ent_user(user: dict = Depends(get_current_user)) -> dict: """企业管理能力端点登录态(能力校验在各公司维度做 is_admin)。""" return user @@ -186,6 +209,8 @@ async def ent_add_member( await sync_user_affiliation(db, target_id) await write_audit(db, action="ent.member_bind", resource="company_member", resource_id=target_id, detail=f"company={cid}", user=user, request=request) + # 成员变更后同步企业群(fire-and-forget,IM 不可用不阻断主流程) + _fire_enterprise_sync(cid) return r @@ -205,6 +230,8 @@ async def ent_remove_member( await sync_user_affiliation(db, uid) await write_audit(db, action="ent.member_remove", resource="company_member", resource_id=uid, detail=f"company={cid}", user=user, request=request) + # 成员变更后同步企业群(fire-and-forget,IM 不可用不阻断主流程) + _fire_enterprise_sync(cid) return r diff --git a/app/api/routers/rbac_opc.py b/app/api/routers/rbac_opc.py index c21142e..87cefc1 100644 --- a/app/api/routers/rbac_opc.py +++ b/app/api/routers/rbac_opc.py @@ -354,7 +354,7 @@ async def opc_affairs( @router.get("/compute/base", summary="算力基础 API 地址") async def opc_compute_base( request: Request, - user: dict = Depends(require_roles("opc_member")), + user: dict = Depends(require_roles("opc_member", "operator")), ): """用户接入地址:经 server-core /v1 中继(OpenAI 兼容)。优先用公网配置,回落请求 host。""" base = (config.COMPUTE_PUBLIC_BASE or str(request.base_url)).rstrip("/") @@ -374,17 +374,17 @@ async def opc_compute_base( @router.get("/compute/models", summary="可用模型(算力中心,admin 新增)") -async def opc_compute_models(_u: dict = Depends(require_roles("opc_member"))): +async def opc_compute_models(_u: dict = Depends(require_roles("opc_member", "operator"))): return {"items": await compute_catalog.models()} @router.get("/compute/prices", summary="模型价目(admin 模型)") -async def opc_compute_prices(_u: dict = Depends(require_roles("opc_member"))): +async def opc_compute_prices(_u: dict = Depends(require_roles("opc_member", "operator"))): return {"items": await compute_catalog.prices()} @router.get("/compute/usage", summary="我的用量(本月按模型)") -async def opc_compute_usage(user: dict = Depends(require_roles("opc_member"))): +async def opc_compute_usage(user: dict = Depends(require_roles("opc_member", "operator"))): try: data = await compute_client.user_usage(user.get("username")) except compute_client.ComputeError as exc: @@ -396,7 +396,7 @@ async def opc_compute_usage(user: dict = Depends(require_roles("opc_member"))): async def opc_compute_usage_logs( page: int = 1, page_size: int = 20, - user: dict = Depends(require_roles("opc_member")), + user: dict = Depends(require_roles("opc_member", "operator")), ): """消耗流水明细,支持分页。""" try: @@ -444,7 +444,7 @@ async def opc_compute_usage_logs( @router.get("/compute/usage/trend", summary="消耗趋势统计(近7天/模型分布/时段分布)") async def opc_compute_usage_trend( days: int = 7, - user: dict = Depends(require_roles("opc_member")), + user: dict = Depends(require_roles("opc_member", "operator")), ): """消耗趋势统计:近N天消耗趋势、模型分布、时段分布。""" from collections import defaultdict @@ -537,7 +537,7 @@ async def opc_compute_usage_trend( @router.get("/compute/balance", summary="我的算力余额") -async def opc_compute_balance(user: dict = Depends(require_roles("opc_member"))): +async def opc_compute_balance(user: dict = Depends(require_roles("opc_member", "operator"))): try: data = await compute_client.user_balance(user.get("username")) except compute_client.ComputeError as exc: @@ -546,7 +546,7 @@ async def opc_compute_balance(user: dict = Depends(require_roles("opc_member"))) @router.get("/compute/tokens", summary="我的算力令牌") -async def opc_compute_tokens(user: dict = Depends(require_roles("opc_member"))): +async def opc_compute_tokens(user: dict = Depends(require_roles("opc_member", "operator"))): """列出当前用户的引擎消费令牌。key 统一补 sk- 前缀(OpenAI 客户端口径)。""" try: data = await compute_client.list_user_tokens(user.get("username")) @@ -574,7 +574,7 @@ class OpcTokenCreateRequest(BaseModel): @router.post("/compute/tokens", summary="创建算力令牌") async def opc_compute_tokens_create( body: OpcTokenCreateRequest | None = None, - user: dict = Depends(require_roles("opc_member")), + user: dict = Depends(require_roles("opc_member", "operator")), ): """为当前用户签发一枚消费令牌(供应商自动对接 PineAgents)。""" try: @@ -591,7 +591,7 @@ async def opc_compute_tokens_create( @router.delete("/compute/tokens/{token_id}", summary="删除算力令牌") -async def opc_compute_tokens_delete(token_id: int, user: dict = Depends(require_roles("opc_member"))): +async def opc_compute_tokens_delete(token_id: int, user: dict = Depends(require_roles("opc_member", "operator"))): # 平台中继令牌不允许删除(桌面端智能体依赖此令牌调用模型) try: existing = await compute_client.list_user_tokens(user.get("username")) diff --git a/app/api/routers/rbac_org.py b/app/api/routers/rbac_org.py index a5754b2..4aa8a1c 100644 --- a/app/api/routers/rbac_org.py +++ b/app/api/routers/rbac_org.py @@ -2,6 +2,8 @@ """组织/租户端点:企业 / 载体 / 服务商(按角色 + 组织作用域)。""" from __future__ import annotations +import logging + from fastapi import APIRouter, Depends, HTTPException from pydantic import BaseModel @@ -9,6 +11,9 @@ from ..dependencies import get_db, get_user_organizations, get_user_primary_org from ..schemas.org import OrgMemberAdd from ...rbac import require_roles from ...infrastructure.repositories import Database +from ...park import tenants as tnt + +log = logging.getLogger("rbac.org") router = APIRouter(tags=["org"]) @@ -22,6 +27,20 @@ async def org_me( orgs = await get_user_organizations(db, user["id"]) region = await db.regions.get(user["region_id"]) if user.get("region_id") else None primary_org = orgs[0] if orgs else None + + # 企业成员关系(company_members 多对多;前端据此渲染企业切换/企业权限) + company_memberships = [] + for m in await db.company_members.list_by_user(user["id"]): + co = await tnt.get_company(m["company_id"]) + company_memberships.append({ + "company_id": m["company_id"], + "company_name": (co or {}).get("name", ""), + "is_admin": bool(m.get("is_admin")), + "member_type": m.get("member_type", ""), + "status": m.get("status", ""), + "joined_at": m.get("created_at", ""), + }) + return { "user_id": user["id"], "role": user["role"], @@ -29,6 +48,7 @@ async def org_me( "org": primary_org, "organizations": orgs, "region": region, + "company_memberships": company_memberships, } @@ -50,11 +70,24 @@ async def carrier_enterprises( async def org_members( org_id: str, db: Database = Depends(get_db), - user: dict = Depends(require_roles("carrier", "operator")), + user: dict = Depends(require_roles("opc_member", "carrier", "operator")), ): - # 机构管理员本人 或 平台运营 可查看 - if "operator" not in (user.get("capabilities") or []) and not await db.org_members.is_admin(org_id, user["id"]): - raise HTTPException(status_code=403, detail="Forbidden: not org admin") + # 叠加制权限:operator 全量;本 org 的载体机构管理员;本企业的企业成员(org_id 亦可能是 company_id) + caps = user.get("capabilities") or [] + if "operator" in caps: + pass + elif await db.org_members.is_admin(org_id, user["id"]): + pass + elif await db.company_members.get(user["id"], org_id): + pass + else: + raise HTTPException(status_code=403, detail="无权查看该组织成员") + + # org_id 可能是 organization_members.org_id(载体机构),也可能被前端当作 company_id 传入: + # 优先返回企业成员(命中企业时),否则回退到机构成员。 + company_members = await tnt.company_members(org_id) + if company_members: + return {"items": company_members} return {"items": await db.org_members.list_for_org(org_id)} diff --git a/app/api/routers/rbac_permissions.py b/app/api/routers/rbac_permissions.py new file mode 100644 index 0000000..1fd8ddb --- /dev/null +++ b/app/api/routers/rbac_permissions.py @@ -0,0 +1,78 @@ +# -*- coding: utf-8 -*- +"""统一权限聚合端点:前端一次性拉取当前用户的完整权限与组织归属。""" +from __future__ import annotations + +import logging + +from fastapi import APIRouter, Depends +from sqlalchemy import select + +from ..dependencies import get_db, get_current_user +from ...infrastructure.models import Organization, OrganizationMember, ParkTenant +from ...infrastructure.repositories import Database +from ...park import tenants as tnt + +log = logging.getLogger("rbac.permissions") + +router = APIRouter(tags=["permissions"]) + + +@router.get("/permissions/me", summary="当前用户完整权限聚合") +async def my_permissions( + db: Database = Depends(get_db), + user: dict = Depends(get_current_user), +): + """聚合角色/能力/权限码 + 园区/企业/机构三类归属,供前端一次性渲染。""" + uid = user["id"] + caps = user.get("capabilities") or [] + + # 园区管理员:park_tenants.operator_user_id 指向当前用户 + park_admin_rows = (await db.session.execute( + select(ParkTenant.id, ParkTenant.name).where(ParkTenant.operator_user_id == uid) + )).all() + park_admin_of = [{"id": r[0], "name": r[1] or ""} for r in park_admin_rows] + + # 企业成员/管理员:company_members WHERE user_id = uid(补企业名称) + company_admin_of: list[dict] = [] + company_member_of: list[dict] = [] + for m in await db.company_members.list_by_user(uid): + co = await tnt.get_company(m["company_id"]) + item = { + "company_id": m["company_id"], + "company_name": (co or {}).get("name", ""), + "is_admin": bool(m.get("is_admin")), + "member_type": m.get("member_type", ""), + "status": m.get("status", ""), + } + company_member_of.append(item) + if item["is_admin"]: + company_admin_of.append(item) + + # 机构(载体)成员:organization_members WHERE user_id = uid(补机构名称/类型) + org_member_of: list[dict] = [] + for m in (await db.session.scalars( + select(OrganizationMember).where(OrganizationMember.user_id == uid))).all(): + org = await db.session.get(Organization, m.org_id) + org_member_of.append({ + "org_id": m.org_id, + "name": org.name if org else "", + "type": org.type if org else "", + "role": m.role, + "is_admin": bool(m.is_admin), + "status": m.status, + }) + + return { + "user_id": uid, + "role": user.get("role"), + "sub_role": user.get("sub_role"), + "capabilities": caps, + "permissions": user.get("permissions", []), + "is_operator": user.get("role") == "operator" or "operator" in caps, + "is_carrier": user.get("role") == "carrier" or "carrier" in caps, + "is_enterprise_admin": "enterprise" in caps, + "park_admin_of": park_admin_of, + "company_admin_of": company_admin_of, + "company_member_of": company_member_of, + "org_member_of": org_member_of, + } diff --git a/app/api/routers/rbac_public.py b/app/api/routers/rbac_public.py index c531631..4e704ce 100644 --- a/app/api/routers/rbac_public.py +++ b/app/api/routers/rbac_public.py @@ -24,3 +24,64 @@ async def list_site_config( """ all_config = await db.config.all() return [item for item in all_config if str(item.get("key", "")).startswith("site.")] + + +@router.get("/startup-resources", summary="启动页轮播资源列表(开放接口,无需权限)") +async def get_startup_resources( + db: Database = Depends(get_db), +): + """返回启动页全屏轮播图的资源列表和配置。 + + 资源列表从配置表 ``site.startup_resources`` 读取;如果未配置,返回兜底默认资源。 + 支持图片和视频,轮播时间、顺序、效果等由接口一次性返回。 + + 返回格式: + { + "resources": [ + { + "url": "https://.../image.jpg", + "type": "image", // image 或 video + "duration": 5000, // 该资源展示时长(毫秒) + "effect": "fade", // 轮播效果:fade/slide/3d/ink/zoom/blur + "order": 1 // 播放顺序 + } + ], + "settings": { + "loop": true, // 是否循环播放 + "shuffle": false, // 是否随机顺序 + "transition_duration": 1200 // 过渡动画时长(毫秒) + } + } + """ + # 尝试从配置表读取自定义资源列表 + try: + config_item = await db.config.get("site.startup_resources") + if config_item and config_item.get("value"): + import json + value = config_item["value"] + if isinstance(value, str): + parsed = json.loads(value) + else: + parsed = value + if isinstance(parsed, dict) and "resources" in parsed: + return parsed + except Exception: + pass + + # 兜底默认资源:本地 bg.mp4 视频循环播放 + return { + "resources": [ + { + "url": "/bg.mp4", + "type": "video", + "duration": 0, # 视频播放完自动切换,0 表示不限制 + "effect": "fade", + "order": 1, + } + ], + "settings": { + "loop": True, + "shuffle": False, + "transition_duration": 1500, + }, + } diff --git a/app/im/client.py b/app/im/client.py index 6c6846e..9ec4087 100644 --- a/app/im/client.py +++ b/app/im/client.py @@ -117,6 +117,19 @@ async def sync_park_group(park_id: str, park_name: str = "") -> bool: return False +async def sync_enterprise_group(company_id: str, company_name: str = "") -> bool: + """确保企业群存在/与最新成员一致。企业成员变更后调用;失败仅记录,不抛给主流程。""" + try: + await _request( + method="POST", path="/internal/enterprise-sync", + headers=_internal_headers(), + json_body={"company_id": company_id, "company_name": company_name}, + ) + return True + except IMError: + return False + + async def internal_mqtt_credentials(user_id: str) -> dict: """登录时生成/更新用户 MQTT 凭证并同步 EMQX(随登录响应下发)。 diff --git a/app/im/router.py b/app/im/router.py index b034a78..35230bb 100644 --- a/app/im/router.py +++ b/app/im/router.py @@ -59,6 +59,18 @@ async def internal_park_sync(request: Request): return {"ok": True} +@router.post("/internal/enterprise-sync") +async def internal_enterprise_sync(request: Request): + body = await _json_body(request) + company_id = (body or {}).get("company_id", "") + if not company_id: + raise HTTPException(status_code=400, detail="company_id 必填") + ok = await client.sync_enterprise_group(company_id, (body or {}).get("company_name", "")) + if not ok: + raise HTTPException(status_code=502, detail="IM 服务不可用") + return {"ok": True} + + # ── REST 转发 ──────────────────────────────────────────────── @router.api_route("/{path:path}", methods=["GET", "POST", "PUT", "PATCH", "DELETE"]) diff --git a/app/main.py b/app/main.py index 5e59168..c2cfab1 100644 --- a/app/main.py +++ b/app/main.py @@ -26,6 +26,7 @@ from app.api.routers import rbac_operator as rbac_operator_router from app.api.routers import rbac_agents as rbac_agents_router from app.api.routers import rbac_public as rbac_public_router from app.api.routers import rbac_org as rbac_org_router +from app.api.routers import rbac_permissions as rbac_permissions_router from app.api.routers import rbac_portals as rbac_portals_router from app.api.routers import rbac_certifications as rbac_certifications_router from app.api.routers import rbac_training as rbac_training_router @@ -107,6 +108,7 @@ app.include_router(rbac_operator_router.router) app.include_router(rbac_agents_router.router) app.include_router(rbac_public_router.router) app.include_router(rbac_org_router.router) +app.include_router(rbac_permissions_router.router) app.include_router(rbac_portals_router.router) app.include_router(rbac_hall_router.router) app.include_router(rbac_service_market_router.router) diff --git a/app/services/compute_catalog.py b/app/services/compute_catalog.py index 3ffde34..5bac9d4 100644 --- a/app/services/compute_catalog.py +++ b/app/services/compute_catalog.py @@ -55,6 +55,21 @@ def _map(row: dict) -> dict: "author_icon": row.get("author_icon", ""), "supplier": row.get("supplier", ""), "supplier_icon": row.get("supplier_icon", ""), + "supplier_icon_file": row.get("supplier_icon_file", ""), + # 封面与系列(模型展示) + "cover": row.get("cover", ""), + "cover_file": row.get("cover_file", ""), + "series": row.get("series", ""), + "series_icon": row.get("series_icon", ""), + "series_icon_file": row.get("series_icon_file", ""), + "notes": row.get("notes", ""), + # 缓存计费:缓存输入 / 缓存写入 + 区间 + "cache_input_price": row.get("cache_input_price", 0) or 0, + "cache_input_price_min": row.get("cache_input_price_min", 0) or 0, + "cache_input_price_max": row.get("cache_input_price_max", 0) or 0, + "cache_write_price": row.get("cache_write_price", 0) or 0, + "cache_write_price_min": row.get("cache_write_price_min", 0) or 0, + "cache_write_price_max": row.get("cache_write_price_max", 0) or 0, "price_is_range": bool(row.get("price_is_range", False)), "input_price_min": row.get("input_price_min", 0) or 0, "input_price_max": row.get("input_price_max", 0) or 0, diff --git a/nginx/nginx/conf.d/opc.pinesound.cn.conf b/nginx/nginx/conf.d/opc.pinesound.cn.conf index 930bfd0..494ced5 100644 --- a/nginx/nginx/conf.d/opc.pinesound.cn.conf +++ b/nginx/nginx/conf.d/opc.pinesound.cn.conf @@ -9,7 +9,7 @@ # 1. nginx 容器必须加入业务网络,才能解析容器名: # docker network connect opc-network nginx # (若 nginx 由 compose 管理,则在 nginx 服务下声明 networks: [opc-network]) -# 2. 禁止再用公网 IP/宿主机回环反代(此前 47.108.226.213:8090 为已废弃旧服务器, +# 2. 禁止再用公网 IP/宿主机回环反代(此前 192.168.1.3:8090 为已废弃旧服务器, # 请求连到错误主机导致 502 upstream prematurely closed)。 # ============================================================ diff --git a/scripts/.sync_event_name_state.json b/scripts/.sync_event_name_state.json new file mode 100644 index 0000000..bdd39a7 --- /dev/null +++ b/scripts/.sync_event_name_state.json @@ -0,0 +1,5 @@ +{ + "ETL33LOE": { + "last_created_at": "2026-09-11T11:42:47.001029Z" + } +} \ No newline at end of file diff --git a/scripts/import_kunming_park.py b/scripts/import_kunming_park.py index 0f65244..c957fdb 100644 --- a/scripts/import_kunming_park.py +++ b/scripts/import_kunming_park.py @@ -14,7 +14,7 @@ import asyncmy import openpyxl import xlrd -DB = dict(host="47.108.226.213", port=8091, user="opc", +DB = dict(host="192.168.1.3", port=8091, user="opc", password="jjjsgysyujkwjwgb", database="opc") TENANT_ID = "T001" BASE = Path("/Volumes/Pine/mycode/opc/昆明市大学生创业园资料") diff --git a/scripts/recover_from_binlog.py b/scripts/recover_from_binlog.py index ccda6fe..585632d 100644 --- a/scripts/recover_from_binlog.py +++ b/scripts/recover_from_binlog.py @@ -96,7 +96,7 @@ async def main(): print(f"解析到 {len(data)} 个表的删除事件") conn = await asyncmy.connect( - host="47.108.226.213", port=8091, + host="192.168.1.3", port=8091, user="root", password="sjnxhyjashaywiuwhaja", database="opc") total_recovered = 0 diff --git a/scripts/sync_event_name_to_nickname.py b/scripts/sync_event_name_to_nickname.py new file mode 100644 index 0000000..c611f3b --- /dev/null +++ b/scripts/sync_event_name_to_nickname.py @@ -0,0 +1,232 @@ +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- +"""同步活动报名「真实姓名」→ 用户昵称(增量脚本,手动不定时运行)。 + +背景 +---- +活动报名时多数用户未在 bookings.name 填真实姓名(默认占位「微信用户」), +真实姓名存于报名表单 form_data_json 的「姓名」字段(events.form_fields_json +中 label=「姓名」的字段)。本脚本按 bookings.username = users.username +关联用户,把报名真实姓名写入 users.nickname。 + +增量机制 +-------- +- state 文件(默认 scripts/.sync_event_name_state.json)记录每个活动上次 + 同步到的最大 created_at;本次只处理 created_at 大于游标的新增报名。 +- 同一用户多条报名只取最新一条。 +- 幂等:昵称已等于目标姓名的不重复写。 +- --full:忽略游标全量重扫(新增历史报名补同步)。 + +覆盖策略 +-------- +- 默认只填充「空 / 占位昵称」(空串、等于 username、等于『微信用户』)的 + 用户,不覆盖用户已设置的真实昵称。 +- --overwrite:强制把昵称改为报名姓名(覆盖已有昵称)。 + +生产连接 +-------- +优先用命令行参数;否则依次读取以下 .env 的 +PINEAGENTS_DEMO_DATABASE_URL(或 PINEAGENTS_MYSQL_* 分字段): + 1. serverrun/.env(生产部署目录,TenXun-Ubuntu2) + 2. .env(项目根) + +用法 +---- + python scripts/sync_event_name_to_nickname.py --dry-run + python scripts/sync_event_name_to_nickname.py + python scripts/sync_event_name_to_nickname.py --full --overwrite --dry-run +""" +from __future__ import annotations + +import argparse +import asyncio +import json +import os +import re +import sys +from datetime import datetime, timezone +from pathlib import Path + +try: + from dotenv import load_dotenv +except ImportError: # 无 dotenv 时跳过(连接参数必须显式给出) + def load_dotenv(*_a, **_k): # noqa: ARG001 + return False + +try: + import asyncmy +except ImportError: + print("缺少依赖 asyncmy,请先安装:pip install asyncmy", file=sys.stderr) + raise SystemExit(1) + +SCRIPT_DIR = Path(__file__).resolve().parent +PROJECT_ROOT = SCRIPT_DIR.parent +DEFAULT_EVENT_ID = "ETL33LOE" # OPC青年数字电商创训营 +DEFAULT_STATE_FILE = SCRIPT_DIR / ".sync_event_name_state.json" + +PLACEHOLDER_NAMES = {"微信用户", "游客", "未填写", ""} + + +def _load_db_config(args: argparse.Namespace) -> dict: + """解析生产数据库连接:参数 > serverrun/.env > .env。""" + if args.host: + return {"host": args.host, "port": args.port, "user": args.user, + "password": args.password, "db": args.db} + + for env_path in (PROJECT_ROOT / "serverrun" / ".env", PROJECT_ROOT / ".env"): + if not env_path.exists(): + continue + load_dotenv(env_path) + url = os.environ.get("PINEAGENTS_DEMO_DATABASE_URL", "") + m = re.match(r"^\w+\+?\w*://([^:]+):([^@]+)@([^:/]+):(\d+)/([^?]+)", url) + if m: + return {"host": m.group(3), "port": int(m.group(4)), + "user": m.group(1), "password": m.group(2), "db": m.group(5)} + host = os.environ.get("PINEAGENTS_MYSQL_HOST", "") + if host: + return {"host": host, + "port": int(os.environ.get("PINEAGENTS_MYSQL_PORT", "3306")), + "user": os.environ.get("PINEAGENTS_MYSQL_USER", ""), + "password": os.environ.get("PINEAGENTS_MYSQL_PASSWORD", ""), + "db": os.environ.get("PINEAGENTS_MYSQL_DBNAME", "opc")} + break # 只在第一个存在的 .env 中取配置 + + print("未找到生产数据库配置:请用 --host/--port/--user/--password/--db 指定," + "或确认 serverrun/.env / .env 存在。", file=sys.stderr) + raise SystemExit(1) + + +def _parse_iso(s: str) -> datetime: + return datetime.fromisoformat(s.replace("Z", "+00:00")) + + +def _is_placeholder_name(name: str) -> bool: + return name.strip() in PLACEHOLDER_NAMES + + +async def sync(args: argparse.Namespace) -> int: + cfg = _load_db_config(args) + print(f"[连接] 生产库 mysql://{cfg['user']}@{cfg['host']}:{cfg['port']}/{cfg['db']}") + conn = await asyncmy.connect(host=cfg["host"], port=cfg["port"], user=cfg["user"], + password=cfg["password"], db=cfg["db"], charset="utf8mb4") + try: + cur = await conn.cursor() + + # 1. 活动表单:定位「姓名」字段 + await cur.execute("SELECT form_fields_json FROM events WHERE id=%s", (args.event_id,)) + row = await cur.fetchone() + if not row: + print(f"活动不存在: {args.event_id}", file=sys.stderr) + return 2 + fields = json.loads(row[0] or "[]") + name_field = next((f.get("id") for f in fields if f.get("label") == "姓名"), None) + if not name_field: + print(f"活动 {args.event_id} 未定义「姓名」字段", file=sys.stderr) + return 2 + print(f"[字段] 姓名字段: {name_field}") + + # 2. 增量游标 + state = {} + if args.state_file.exists(): + try: + state = json.loads(args.state_file.read_text(encoding="utf-8")) + except (OSError, json.JSONDecodeError): + state = {} + cursor_ts = "" if args.full else (state.get(args.event_id) or {}).get("last_created_at", "") + if cursor_ts: + print(f"[增量] 游标 last_created_at = {cursor_ts}(--full 可全量重扫)") + else: + print("[增量] 无历史游标,全量扫描") + + # 3. 取报名(只取该活动、created_at > 游标) + sql = ("SELECT id, username, name, form_data_json, created_at FROM bookings " + "WHERE event_id=%s") + params = [args.event_id] + if cursor_ts: + sql += " AND created_at > %s" + params.append(cursor_ts) + sql += " ORDER BY created_at ASC" + await cur.execute(sql, params) + rows = await cur.fetchall() + print(f"[扫描] 待处理报名 {len(rows)} 条") + + # 4. 同一用户多条报名只保留最新一条 + latest: dict[str, dict] = {} + for bid, username, name, form_json, created_at in rows: + if not username: + continue + fd = json.loads(form_json or "{}") + real_name = str(fd.get(name_field) or "").strip() + latest[username] = {"booking_id": bid, "name": name or "", + "real_name": real_name, "created_at": created_at} + + # 5. 关联用户并执行 + updated, skipped_no_name, skipped_placeholder, skipped_missing, skipped_same = 0, 0, 0, 0, 0 + max_created = cursor_ts + for username, info in latest.items(): + if info["created_at"] > max_created: + max_created = info["created_at"] + real = info["real_name"] + if not real: + skipped_no_name += 1 + continue + if _is_placeholder_name(real) or real == username or len(real) < 2: + skipped_placeholder += 1 + continue + + await cur.execute("SELECT username, nickname FROM users WHERE username=%s", (username,)) + u = await cur.fetchone() + if not u: + skipped_missing += 1 + print(f" [跳过] 用户不存在 {username}") + continue + _u, nickname = u[0], u[1] or "" + if nickname == real: + skipped_same += 1 + continue + if not args.overwrite and not ( + not nickname.strip() + or nickname.strip() == username + or _is_placeholder_name(nickname.strip()) + ): + skipped_same += 1 # 已有真实昵称且未开 --overwrite,不动 + print(f" [保留] {username} 已有昵称「{nickname}」,不覆盖") + continue + print(f" [{'DRY' if args.dry_run else 'SET '}] {username} 昵称: " + f"{nickname or '(空)'} -> {real} (报名 {info['booking_id']})") + if not args.dry_run: + await cur.execute("UPDATE users SET nickname=%s WHERE username=%s", + (real, username)) + updated += 1 + + if not args.dry_run: + await conn.commit() + # 更新游标(幂等,重复跑安全) + state.setdefault(args.event_id, {})["last_created_at"] = max_created + args.state_file.parent.mkdir(parents=True, exist_ok=True) + args.state_file.write_text( + json.dumps(state, ensure_ascii=False, indent=2), encoding="utf-8") + print(f"[游标] 已更新 -> {max_created}") + print(f"[结果] {'(dry-run 未写库)' if args.dry_run else '已写库'}" + f" 更新 {updated};跳过:无姓名 {skipped_no_name}、占位/无效 {skipped_placeholder}、" + f"用户不存在 {skipped_missing}、无需变更 {skipped_same}") + return 0 + finally: + conn.close() + + +if __name__ == "__main__": + parser = argparse.ArgumentParser(description="同步活动报名真实姓名到用户昵称(增量)") + parser.add_argument("--event-id", default=DEFAULT_EVENT_ID, help="活动 ID") + parser.add_argument("--dry-run", action="store_true", help="只预览不写库") + parser.add_argument("--full", action="store_true", help="忽略增量游标全量重扫") + parser.add_argument("--overwrite", action="store_true", help="覆盖用户已有昵称") + parser.add_argument("--state-file", type=Path, default=DEFAULT_STATE_FILE, + help="增量游标文件路径") + parser.add_argument("--host", help="MySQL 主机(默认读 .env)") + parser.add_argument("--port", type=int, default=0, help="MySQL 端口") + parser.add_argument("--user", help="MySQL 用户") + parser.add_argument("--password", help="MySQL 密码") + parser.add_argument("--db", default="opc", help="MySQL 库名") + args = parser.parse_args() + raise SystemExit(asyncio.run(sync(args))) diff --git a/serverdata/backup/opc_full_20260908_175552.sql b/serverdata/backup/opc_full_20260908_175552.sql index a0c0ef5..7ccca53 100644 --- a/serverdata/backup/opc_full_20260908_175552.sql +++ b/serverdata/backup/opc_full_20260908_175552.sql @@ -1,6 +1,6 @@ -- MySQL dump 10.13 Distrib 26.7.0, for macos26.6 (arm64) -- --- Host: 47.108.226.213 Database: opc +-- Host: 192.168.1.3 Database: opc -- ------------------------------------------------------ -- Server version 8.4.11 diff --git a/serverrun/.env.example b/serverrun/.env.example index 1d52257..2d63383 100644 --- a/serverrun/.env.example +++ b/serverrun/.env.example @@ -6,7 +6,7 @@ # 默认值 = 本地 Docker 基础设施(serverrun/mysql、redis、mqtt 子目录编排): # 先 docker network create opc-network,再启动基础设施,业务容器经 # 容器名(opc-mysql / opc-redis / opc-emqx)+ 内部端口互访。 -# 生产环境:替换为外部服务器地址(每段注释附 47.108.226.213 示例)。 +# 生产环境:替换为外部服务器地址(每段注释附 192.168.1.3 示例)。 # # 容器间通过 compose 服务名互访(不再是 127.0.0.1 回环): # opc-server → http://opc-compute:3000 / http://opc-im:8101 @@ -18,7 +18,7 @@ # ───────────────────────────────────────────── # MySQL(本地 Docker:opc-mysql 容器内部端口 3306;宿主机访问用 localhost:8091) -# 生产:PINEAGENTS_MYSQL_HOST=47.108.226.213 / PINEAGENTS_MYSQL_PORT=8091 +# 生产:PINEAGENTS_MYSQL_HOST=192.168.1.3 / PINEAGENTS_MYSQL_PORT=8091 PINEAGENTS_MYSQL_HOST=opc-mysql PINEAGENTS_MYSQL_PORT=3306 PINEAGENTS_MYSQL_USER=opc @@ -27,7 +27,7 @@ PINEAGENTS_MYSQL_DB=opc PINEAGENTS_DEMO_DATABASE_URL=mysql+asyncmy://opc:jjjsgysyujkwjwgb@opc-mysql:3306/opc?charset=utf8mb4 # Redis(本地 Docker:opc-redis,密码见 redis/docker-compose.yml 的 requirepass) -# 生产:PINEAGENTS_REDIS_HOST=47.108.226.213 / PINEAGENTS_REDIS_PORT=6379 +# 生产:PINEAGENTS_REDIS_HOST=192.168.1.3 / PINEAGENTS_REDIS_PORT=6379 PINEAGENTS_REDIS_HOST=opc-redis PINEAGENTS_REDIS_PORT=6379 PINEAGENTS_REDIS_PASSWORD=redis_pass @@ -66,14 +66,14 @@ IM_PORT=8101 # 共享 opc MySQL(本地 Docker) IM_DATABASE_URL=mysql+asyncmy://opc:jjjsgysyujkwjwgb@opc-mysql:3306/opc?charset=utf8mb4 -# MQTT / EMQX 总线(本地 Docker:opc-emqx 容器内部 1883;生产 MQTT_BROKER_HOST=47.108.226.213 / MQTT_BROKER_PORT=8095) +# MQTT / EMQX 总线(本地 Docker:opc-emqx 容器内部 1883;生产 MQTT_BROKER_HOST=192.168.1.3 / MQTT_BROKER_PORT=8095) IM_MQTT_ENABLED=true MQTT_BROKER_HOST=opc-emqx MQTT_BROKER_PORT=1883 MQTT_USERNAME=opc MQTT_PASSWORD=Opc123! -# EMQX 管理 API(本地 http://opc-emqx:18083;生产 IM_EMQX_API_URL=https://47.108.226.213:8096) +# EMQX 管理 API(本地 http://opc-emqx:18083;生产 IM_EMQX_API_URL=https://192.168.1.3:8096) IM_EMQX_API_URL=http://opc-emqx:18083 IM_EMQX_API_KEY= IM_EMQX_API_SECRET= @@ -86,7 +86,7 @@ IM_INTERNAL_TOKEN= # 智能客服 LLM(可选) IM_CS_LLM_ENABLED=false -# 桌面端 MQTT over WebSocket 直连入口(本地:宿主 localhost:8083;生产 47.108.226.213:8097) +# 桌面端 MQTT over WebSocket 直连入口(本地:宿主 localhost:8083;生产 192.168.1.3:8097) IM_MQTT_PUBLIC_HOST=localhost IM_MQTT_PUBLIC_WS_PORT=8083