"""数据中枢控制台端到端自测:主站与中枢真的对话一遍,不是 mock。 跑法:python tools/verify_datahub_console.py 覆盖邀请码一次性注册、管理员门禁、凭证掩码、模型池与会员桥接读写、服务令牌校验。 两个服务都起在临时端口 + 临时数据目录,跑完自动清理,不碰任何现网数据。 """ from __future__ import annotations import http.cookies import json import os import shutil import socket import sqlite3 import sys import tempfile import threading import time import urllib.error import urllib.request from pathlib import Path ROOT = Path(__file__).resolve().parents[1] # 从 tools/ 下运行,仓库根不在 sys.path 上;主站包按仓库根导入。 if str(ROOT) not in sys.path: sys.path.insert(0, str(ROOT)) HUB_TOKEN = "smoke-hub-admin-token-0123456789abcdef" PASSWORD = "SmokePass123" FAILURES: list[str] = [] def free_port() -> int: with socket.socket() as sock: sock.bind(("127.0.0.1", 0)) return int(sock.getsockname()[1]) def check(label: str, ok: bool, detail: str = "") -> None: print(f" {'PASS' if ok else 'FAIL'} {label}{(' — ' + detail) if detail else ''}") if not ok: FAILURES.append(label) def request(url: str, payload=None, method="GET", headers=None, cookie="") -> tuple[int, dict, str]: data = json.dumps(payload).encode("utf-8") if payload is not None else None req = urllib.request.Request(url, data=data, method=method) req.add_header("Content-Type", "application/json; charset=utf-8") for key, value in (headers or {}).items(): req.add_header(key, value) if cookie: req.add_header("Cookie", cookie) try: with urllib.request.urlopen(req, timeout=15) as response: raw = response.read().decode("utf-8") set_cookie = response.headers.get("Set-Cookie") or "" return response.status, _parse(raw), set_cookie except urllib.error.HTTPError as exc: return exc.code, _parse(exc.read().decode("utf-8")), exc.headers.get("Set-Cookie") or "" def _parse(raw: str) -> dict: try: parsed = json.loads(raw) except json.JSONDecodeError: return {"_raw": raw[:200]} return parsed if isinstance(parsed, dict) else {"_list": parsed} def session_cookie(header: str) -> str: jar = http.cookies.SimpleCookie() jar.load(header) morsel = jar.get("xiaobai_session") return f"xiaobai_session={morsel.value}" if morsel else "" def wait_for(url: str, seconds: float = 20.0) -> bool: deadline = time.time() + seconds while time.time() < deadline: try: urllib.request.urlopen(url, timeout=2) return True except urllib.error.HTTPError: return True except OSError: time.sleep(0.25) return False REVIEW_DB = ROOT / "data" / "review.db" ENV_FILE = ROOT / ".env" class RepoSandbox: """主站的库路径和 .env 都写死在仓库里,跑之前挪开、跑完原样放回。""" def __enter__(self) -> "RepoSandbox": self.stash = Path(tempfile.mkdtemp(prefix="hel560-stash-")) for path in (REVIEW_DB, ENV_FILE): if path.exists(): shutil.copy2(path, self.stash / path.name) if REVIEW_DB.exists(): REVIEW_DB.unlink() # 自测需要一个空库来验证"首个账号免邀请码" return self def __exit__(self, *exc_info: object) -> None: for path in (REVIEW_DB, ENV_FILE): saved = self.stash / path.name if saved.exists(): shutil.copy2(saved, path) elif path.exists(): path.unlink() for extra in REVIEW_DB.parent.glob("review.db-*"): extra.unlink() shutil.rmtree(self.stash, ignore_errors=True) def start_review(workdir: Path, port: int) -> None: os.environ["HUB_ADMIN_TOKEN"] = HUB_TOKEN from backend.application import RequestHandler, SERVICE # noqa: F401 from http.server import ThreadingHTTPServer server = ThreadingHTTPServer(("127.0.0.1", port), RequestHandler) threading.Thread(target=server.serve_forever, daemon=True).start() def start_hub(workdir: Path, port: int, review_port: int) -> None: sys.path.insert(0, str(ROOT / "xiaobai-datahub")) from cryptography.fernet import Fernet os.environ.update( { "DATAHUB_DB_PATH": str(workdir / "hub.db"), "DATAHUB_BACKUP_DIR": str(workdir / "backups"), "DATAHUB_ENCRYPTION_KEY": Fernet.generate_key().decode(), "DATAHUB_TOKEN": "smoke-datahub-token", "HUB_ADMIN_TOKEN": HUB_TOKEN, "REVIEW_BASE_URL": f"http://127.0.0.1:{review_port}", "REVIEW_PUBLIC_URL": f"http://127.0.0.1:{review_port}", "DATAHUB_SCHEDULER_ENABLED": "0", } ) from datahub.httpapp import make_handler from datahub.hub import Hub from datahub.settings import load_settings from http.server import ThreadingHTTPServer hub = Hub(load_settings(os.environ)) server = ThreadingHTTPServer(("127.0.0.1", port), make_handler(hub)) threading.Thread(target=server.serve_forever, daemon=True).start() def main() -> int: if any(argument in {"-h", "--help"} for argument in sys.argv[1:]): print(__doc__.strip()) return 0 workdir = Path(tempfile.mkdtemp(prefix="hel560-smoke-")) review_port, hub_port = free_port(), free_port() review = f"http://127.0.0.1:{review_port}" hub = f"http://127.0.0.1:{hub_port}" with RepoSandbox(): start_review(workdir, review_port) if not wait_for(f"{review}/api/session"): print("主站没起来") return 1 start_hub(workdir, hub_port, review_port) if not wait_for(f"{hub}/livez"): print("数据中枢没起来") return 1 print("\n[1] 首个账号免邀请码,之后注册强制邀请码") status, body, cookie_header = request(f"{review}/api/auth/register", {"username": "boss", "password": PASSWORD}, "POST") check("首个账号可直接注册(自动成为管理员)", status == 201, f"{status} {body.get('error', '')}") admin_cookie = session_cookie(cookie_header) status, body, _ = request(f"{review}/api/auth/register", {"username": "nobody", "password": PASSWORD}, "POST") check("第二个账号没邀请码被拒", status >= 400 and "邀请码" in str(body.get("error", "")), f"{status} {body.get('error', '')}") print("\n[2] 未登录 / 非管理员进不了数据中枢") status, body, _ = request(f"{hub}/admin/api/session") check("未登录访问控制台返回 401 并给出主站登录地址", status == 401 and "/login/" in str(body.get("login_url", "")), f"{status} {body}") print("\n[3] 主站管理员会话直接进控制台(跨端口共享 cookie)") status, session, _ = request(f"{hub}/admin/api/session", cookie=admin_cookie) check("带主站会话访问控制台返回 200", status == 200, f"{status} {session}") check("控制台回显主站用户名", session.get("username") == "boss", str(session.get("username"))) csrf = str(session.get("csrf") or "") check("下发了 CSRF 令牌", len(csrf) >= 32, csrf[:12]) write_headers = {"X-CSRF-Token": csrf} print("\n[4] 写接口必须带 CSRF") status, body, _ = request(f"{hub}/admin/api/invites/create", {"count": 1}, "POST", cookie=admin_cookie) check("缺 CSRF 的写请求被拒", status == 401, f"{status} {body}") print("\n[5] 控制台生成邀请码 → 注册消耗一次 → 二次使用失败") status, created, _ = request(f"{hub}/admin/api/invites/create", {"count": 2}, "POST", write_headers, admin_cookie) check("控制台生成邀请码成功", status == 200 and len(created.get("created") or []) == 2, f"{status} {created.get('error', '')}") codes = [item["code"] for item in created.get("created") or []] check("列表只给掩码,不回明文", all("•" in row["code_masked"] for row in created.get("codes") or [])) status, body, member_cookie_header = request( f"{review}/api/auth/register", {"username": "xiaochen", "password": PASSWORD, "invite_code": codes[0]}, "POST" ) check("凭邀请码注册成功", status == 201, f"{status} {body.get('error', '')}") member_cookie = session_cookie(member_cookie_header) status, body, _ = request( f"{review}/api/auth/register", {"username": "again", "password": PASSWORD, "invite_code": codes[0]}, "POST" ) check("同一邀请码第二次注册被拒", status >= 400, f"{status} {body.get('error', '')}") print("\n[6] 作废后的邀请码不能注册") handles = {row["code_masked"][:7]: row["code_id"] for row in created.get("codes") or []} target = handles.get(codes[1][:7]) status, body, _ = request(f"{hub}/admin/api/invites/revoke", {"code_id": target}, "POST", write_headers, admin_cookie) check("控制台作废未使用的邀请码", status == 200, f"{status} {body.get('error', '')}") status, body, _ = request( f"{review}/api/auth/register", {"username": "revoked", "password": PASSWORD, "invite_code": codes[1]}, "POST" ) check("已作废邀请码无法注册", status >= 400, f"{status} {body.get('error', '')}") print("\n[7] 普通会员账号进不了控制台") status, body, _ = request(f"{hub}/admin/api/session", cookie=member_cookie) check("非管理员访问控制台返回 403", status == 403, f"{status} {body}") status, body, _ = request(f"{hub}/admin/api/members", cookie=member_cookie) check("非管理员读会员接口同样 403", status == 403, f"{status} {body}") print("\n[8] 数据源凭证在线写入 + 掩码回显") status, body, _ = request( f"{hub}/admin/api/credentials/tushare", {"tushare_token": "tok-abcdefgh1234"}, "POST", write_headers, admin_cookie ) check("Tushare Token 保存成功", status == 200, f"{status} {body.get('error', '')}") status, sources, _ = request(f"{hub}/admin/api/sources", cookie=admin_cookie) tushare = next((row for row in sources.get("items") or [] if row.get("provider") == "tushare"), {}) credential = tushare.get("credential") or {} check("数据源卡回显掩码而非明文", credential.get("configured") and "1234" in str(credential.get("last4")), str(credential)) check("接口不回传明文 Token", "tok-abcdefgh1234" not in json.dumps(sources, ensure_ascii=False)) print("\n[9] 模型池 / 会员 / 邀请码三页都能从控制台读到") for label, path in (("模型池", "/admin/api/models"), ("会员", "/admin/api/members"), ("邀请码", "/admin/api/invites")): status, body, _ = request(f"{hub}{path}", cookie=admin_cookie) check(f"{label}接口可读", status == 200, f"{status} {body.get('error', '')}") print("\n[10] 控制台改模型池 → 主站落库") models = [{"id": "smoke-main", "name": "冒烟主模型", "model": "gpt-4o", "base_url": "https://api.openai.com/v1", "api_key": "sk-smoke-key-9911"}] status, body, _ = request( f"{hub}/admin/api/models/save", {"models": models, "primary_model_id": "smoke-main"}, "POST", write_headers, admin_cookie ) check("控制台保存模型池成功", status == 200, f"{status} {body.get('error', '')}") groups = body.get("groups") or [] check("按供应商归组返回", len(groups) == 1 and groups[0]["base_url"] == "https://api.openai.com/v1", str(groups)[:120]) check("供应商显示密钥后四位而非明文", groups and groups[0].get("key_last4") == "9911", str(groups[0].get("key_last4") if groups else "")) check("模型接口不回传明文密钥", "sk-smoke-key-9911" not in json.dumps(body, ensure_ascii=False)) status, status_body, _ = request( f"{review}/api/admin/settings", None, "GET", {"X-Hub-Admin-Token": HUB_TOKEN}, admin_cookie ) status, mainsite, _ = request(f"{review}/api/hub-admin/status", {}, "POST", {"X-Hub-Admin-Token": HUB_TOKEN}) pool = (mainsite.get("llm") or {}).get("models") or [] check("主站确实存下了这个模型", any(m["id"] == "smoke-main" for m in pool), str([m.get("id") for m in pool])) print("\n[11] 会员额度与会员开通经控制台落到主站") status, body, _ = request(f"{hub}/admin/api/members/quota", {"member_daily_limit": 88}, "POST", write_headers, admin_cookie) check("保存会员每日额度成功", status == 200 and (body.get("membership") or {}).get("member_daily_limit") == 88, f"{status} {body.get('membership')}") member_id = next((u["id"] for u in body.get("users") or [] if u["username"] == "xiaochen"), 0) status, body, _ = request( f"{hub}/admin/api/members/save", {"user_id": member_id, "status": "active", "duration": "3_months"}, "POST", write_headers, admin_cookie ) row = next((u for u in body.get("users") or [] if u["id"] == member_id), {}) check("开通 3 个月会员生效", status == 200 and row.get("membership_status") == "active" and row.get("membership_expires_at"), f"{status} {row.get('membership_status')} {row.get('membership_expires_at')}") print("\n[12] 桥接令牌是唯一信任边界") status, body, _ = request(f"{review}/api/hub-admin/status", {}, "POST", {"X-Hub-Admin-Token": "wrong-token"}) check("桥接端点拒绝错误令牌", status == 401, f"{status} {body}") status, body, _ = request(f"{review}/api/hub-admin/status", {}, "POST") check("桥接端点拒绝无令牌", status == 401, f"{status} {body}") status, body, _ = request(f"{review}/api/hub-admin/invites", {}, "POST", {"X-Hub-Admin-Token": HUB_TOKEN}, admin_cookie) check("带正确令牌可读邀请码", status == 200, f"{status} {body.get('error', '')}") print("\n[13] 退出登录会真的销毁主站会话") status, body, _ = request(f"{hub}/admin/api/logout", {}, "POST", write_headers, admin_cookie) check("控制台退出返回主站登录地址", status == 200 and "/login/" in str(body.get("login_url", "")), f"{status} {body}") status, body, _ = request(f"{review}/api/session", cookie=admin_cookie) check("主站会话已失效", not (body.get("authenticated") or body.get("user")), str(body)[:120]) print("\n[14] 并发使用同一邀请码只成功一次") os.environ["HUB_ADMIN_TOKEN"] = HUB_TOKEN with sqlite3.connect(REVIEW_DB) as connection: rows = connection.execute("SELECT status, COUNT(*) FROM invite_codes GROUP BY status").fetchall() counts = dict(rows) check("邀请码状态落库正确(1 已用 / 1 已作废)", counts.get("used") == 1 and counts.get("revoked") == 1, str(counts)) shutil.rmtree(workdir, ignore_errors=True) print("\n" + "=" * 60) if FAILURES: print(f"FAILED {len(FAILURES)} 项:") for item in FAILURES: print(" - " + item) return 1 print("数据中枢控制台端到端自测全部通过") return 0 if __name__ == "__main__": raise SystemExit(main())