#!/usr/bin/env python3 # -*- coding: utf-8 -*- """观澜云端版 · 服务端 /api/qa 代理 (草案, 未部署; 部署包 CLAUDE_CODE_部署提示词.md §安全边界 / §账号与权限 / §B 云端可用性; 交接 2026-09-06). 边界 (逐条对应部署包): - 密钥只在服务端环境变量 DEEPSEEK_API_KEY; 缺失 → 模式 "不可用" + 明确文案, **不伪造回答**, 不落任何文件/日志/回包. - 浏览器不直连 DeepSeek: 本进程绑 127.0.0.1, 前面放 nginx (HTTPS); 上游 401 / 429 / 超时 / 5xx → 回包 ok=false + 脱敏错误码, answer 恒 null. - 鉴权: 服务端会话 (HttpOnly cookie, 随机 token, 服务端过期); 账号 = 拼音显示名 + 临时密码 (PBKDF2, 10 天有效, 可撤销, 建/撤销/登录全部进审计). - 限流 (按用户 + 按 IP, 固定窗) · 请求长度 (body / 问题字数) · 上游超时 · 审计日志 (JSONL: 只记 问题 sha256+长度, 不记原文, 不记密钥) · 错误脱敏 (上游原文只进服务端 stderr 且先去密钥). - 模式标签: 每个回包带 mode ∈ {云端 DeepSeek, 本机 DeepSeek, 演示回答, 不可用}; 演示回答只在 GUANLAN_QA_DEMO=1 且无密钥时出现, 且 answer 前缀明示 "[演示回答]". - 出口脱敏: 模型答案过 src/windscada/deid_public.scrub (与 windscada_ask_serve 同款; 导入失败 = 拒绝启动). - 接地: 问题先在云端脱敏面孔 (outputs/rudong/guanlan/cloud/claims_public.json) 里按关键词取 ≤5 条作 system 上下文, 回包 refs 列 claim_id; 面孔缺失 → 不带引用, refs=[] (不伪造引用). 不做: 用户管理 HTTP 接口 (只走 CLI, 缩小攻击面) · 多轮对话 · 流式. 用法: serve [--port 8090] | user add [--role user|admin] [--days 10] | user revoke | user list | selftest 环境变量 (名, 不含值): DEEPSEEK_API_KEY · DEEPSEEK_BASE_URL (默认 https://api.deepseek.com) · DEEPSEEK_MODEL (默认 deepseek-chat) · GUANLAN_QA_USERS (users.json 路径) · GUANLAN_QA_AUDIT (audit.jsonl 路径) · GUANLAN_QA_DEMO · GUANLAN_QA_INSECURE_COOKIE (=1 时 cookie 不带 Secure, 仅本机 http 测试) 上游调用经 transport 注入 (ask(q, user, transport=...)), 测试用桩模拟 401/429/超时/5xx, 不真调云. """ import argparse, base64, datetime as dt, hashlib, hmac, http.client, http.cookies, json, os, pathlib, re, secrets, socket, sys, threading, time, urllib.parse, uuid from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer ROOT = pathlib.Path(__file__).resolve().parents[1] sys.path.insert(0, str(ROOT / "src")) from windscada import deid_public as DP # noqa: E402 ★出口脱敏不可用 = 拒绝启动 (看着有闸实则没有, 比没闸更坏) MODES = ("云端 DeepSeek", "本机 DeepSeek", "演示回答", "不可用") LIMITS = dict(max_body_bytes=16 * 1024, max_q_chars=2000, rate_user=(10, 300), rate_ip=(30, 300), upstream_timeout_s=60, max_answer_chars=6000, temp_pw_days=10, session_hours=12, max_refs=5) CLOUD_CLAIMS = ROOT / "outputs/rudong/guanlan/cloud/claims_public.json" USERS_PATH = pathlib.Path(os.environ.get("GUANLAN_QA_USERS", str(ROOT / "outputs/rudong/guanlan/cloud/_private/users.json"))) AUDIT_PATH = pathlib.Path(os.environ.get("GUANLAN_QA_AUDIT", str(ROOT / "outputs/rudong/guanlan/cloud/_private/audit.jsonl"))) SYS_PROMPT = ("你是风电场 SCADA/CMS 分析结论的解释助手。只能依据下面给出的『契约结论』回答; 结论里没有的数字不许编, 没有依据就说『契约里没有这条』。" "不得批准检修、不得替代工程师判断、不得补造缺失数据。回答用中文, 引用结论时写 claim_id。") ERR_TEXT = {"no_key": "云端模型未配置 (服务端缺 DEEPSEEK_API_KEY), 问答不可用", "upstream_401": "云端模型凭证无效 (服务端配置问题), 问答不可用", "upstream_429": "云端模型限流, 请稍后再试", "upstream_timeout": "云端模型超时, 本次未得到回答", "upstream_5xx": "云端模型服务异常, 本次未得到回答", "upstream_net": "云端模型不可达, 本次未得到回答", "upstream_bad": "云端模型回包无法解析, 本次未得到回答"} _KEY_RE = re.compile(r"sk-[A-Za-z0-9]{8,}|Bearer\s+\S+") # ── 密钥 / 脱敏 ──────────────────────────────────────────────────── def api_key(): return (os.environ.get("DEEPSEEK_API_KEY") or "").strip() def redact(s): """任何要落日志/回包的字符串先过这里: 去掉密钥形态串与真实密钥值.""" s = str(s); k = api_key() if k: s = s.replace(k, "[REDACTED]") return _KEY_RE.sub("[REDACTED]", s) def mode_now(): if api_key(): return "云端 DeepSeek" return "演示回答" if os.environ.get("GUANLAN_QA_DEMO") == "1" else "不可用" # ── 审计 ────────────────────────────────────────────────────────── _AUDIT_LOCK = threading.Lock() def audit(event, **kw): rec = {"ts": dt.datetime.now().isoformat(timespec="seconds"), "event": event} rec.update({k: (redact(v) if isinstance(v, str) else v) for k, v in kw.items()}) AUDIT_PATH.parent.mkdir(parents=True, exist_ok=True) with _AUDIT_LOCK, AUDIT_PATH.open("a", encoding="utf-8") as f: f.write(json.dumps(rec, ensure_ascii=False) + "\n") return rec # ── 账号 (服务端管理; 拼音显示名 + 临时密码) ────────────────────── PBKDF2_ITERS = 200_000 def _hash_pw(pw, salt): return hashlib.pbkdf2_hmac("sha256", pw.encode("utf-8"), bytes.fromhex(salt), PBKDF2_ITERS).hex() class Users: def __init__(self, path=USERS_PATH): self.path = pathlib.Path(path); self.d = json.loads(self.path.read_text(encoding="utf-8")) if self.path.is_file() else {"users": {}} def _save(self): self.path.parent.mkdir(parents=True, exist_ok=True); self.path.write_text(json.dumps(self.d, ensure_ascii=False, indent=1), encoding="utf-8") try: os.chmod(self.path, 0o600) except OSError: pass def add(self, name, role="user", days=LIMITS["temp_pw_days"], by="cli"): """建/重置账号 → 临时密码 (只返回这一次, 不落盘明文). 最小权限: 默认 user.""" if not re.fullmatch(r"[a-z][a-z0-9_]{1,31}", name): raise ValueError("显示名须为拼音小写 [a-z0-9_], 2~32 字") if role not in ("user", "admin"): raise ValueError("role ∈ user|admin") pw = base64.b32encode(secrets.token_bytes(10)).decode().rstrip("=").lower(); salt = secrets.token_hex(16) self.d["users"][name] = {"role": role, "salt": salt, "hash": _hash_pw(pw, salt), "created": dt.datetime.now().isoformat(timespec="seconds"), "expires": (dt.datetime.now() + dt.timedelta(days=days)).isoformat(timespec="seconds"), "revoked": False} self._save(); audit("user_add", user=name, role=role, days=days, by=by); return pw def revoke(self, name, by="cli"): if name not in self.d["users"]: raise KeyError(name) self.d["users"][name]["revoked"] = True; self._save(); audit("user_revoke", user=name, by=by) def verify(self, name, pw, now=None): """→ (user dict | None, reason ∈ ok|no_user|bad_pw|expired|revoked). 不区分文案给客户端 (统一 401), reason 只进审计.""" u = self.d["users"].get(name or "") if not u: _hash_pw(pw or "", "00" * 16); return None, "no_user" # 等时: 无账号也算一次 hash if not hmac.compare_digest(_hash_pw(pw or "", u["salt"]), u["hash"]): return None, "bad_pw" if u.get("revoked"): return None, "revoked" if (now or dt.datetime.now()) > dt.datetime.fromisoformat(u["expires"]): return None, "expired" return {"name": name, "role": u["role"], "expires": u["expires"]}, "ok" def list(self): return [{"name": n, "role": u["role"], "expires": u["expires"], "revoked": u.get("revoked", False), "created": u["created"]} for n, u in self.d["users"].items()] class Sessions: def __init__(self, hours=LIMITS["session_hours"]): self.s = {}; self.hours = hours; self.lock = threading.Lock() def create(self, user): tok = secrets.token_urlsafe(32) with self.lock: self.s[tok] = dict(user, exp=time.time() + self.hours * 3600) return tok def get(self, tok): with self.lock: u = self.s.get(tok or "") if u and (u["exp"] < time.time() or dt.datetime.now() > dt.datetime.fromisoformat(u["expires"])): self.s.pop(tok, None); return None return u def drop(self, tok): with self.lock: self.s.pop(tok or "", None) class RateLimiter: """固定窗: key 在 window 秒内最多 n 次 → (ok, retry_after_s).""" def __init__(self): self.w = {}; self.lock = threading.Lock() def allow(self, key, n, window, now=None): now = now if now is not None else time.time() with self.lock: t0, c = self.w.get(key, (now, 0)) if now - t0 >= window: t0, c = now, 0 if c >= n: return False, int(window - (now - t0)) + 1 self.w[key] = (t0, c + 1); return True, 0 # ── 接地引用 (云端脱敏面孔) ───────────────────────────────────────── _CLAIMS = {"key": None, "rows": []} def claims_public(): if not CLOUD_CLAIMS.is_file(): return [] st = CLOUD_CLAIMS.stat(); key = (st.st_mtime_ns, st.st_size) if _CLAIMS["key"] != key: _CLAIMS.update(key=key, rows=json.loads(CLOUD_CLAIMS.read_text(encoding="utf-8")).get("claims", [])) return _CLAIMS["rows"] def pick_refs(q, k=LIMITS["max_refs"]): """关键词重叠取 top-k 契约条 (只用脱敏面孔; 面孔缺 → []). 不是检索系统, 只是把回答钉在契约上.""" toks = {t for t in re.findall(r"[一-鿿]{2,}|[A-Za-z][A-Za-z0-9\-]{2,}", q or "")} grams = {t[i:i + 2] for t in toks for i in range(len(t) - 1)} | toks if not grams: return [] scored = [] for c in claims_public(): text = (c.get("title_public") or "") + " " + (c.get("missing_evidence") or "") sc = sum(1 for g in grams if g in text) + (2 if c.get("claim_id", "").lower() in (q or "").lower() else 0) if sc: scored.append((sc, c)) scored.sort(key=lambda x: (-x[0], x[1]["claim_id"])) return [c for _, c in scored[:k]] def build_messages(q, refs): ctx = "\n".join(f"- {c['claim_id']} [{c['verdict']}] {c['title_public'][:200]}" + (f" | 缺失证据: {c['missing_evidence'][:120]}" if c.get("missing_evidence") else "") for c in refs) or "(契约里没有与本问题相关的结论)" return [{"role": "system", "content": SYS_PROMPT + "\n\n契约结论:\n" + ctx}, {"role": "user", "content": q}] # ── 上游 ────────────────────────────────────────────────────────── class UpstreamError(Exception): def __init__(self, code, detail=""): super().__init__(code); self.code = code; self.detail = redact(detail)[:300] def deepseek_transport(messages, timeout=LIMITS["upstream_timeout_s"]): """真上游 (只此一处发网络请求). → 答案文本; 失败 raise UpstreamError(code ∈ upstream_401|upstream_429|upstream_timeout|upstream_5xx|upstream_net|upstream_bad).""" key = api_key() if not key: raise UpstreamError("no_key") base = urllib.parse.urlparse(os.environ.get("DEEPSEEK_BASE_URL") or "https://api.deepseek.com") body = json.dumps({"model": os.environ.get("DEEPSEEK_MODEL") or "deepseek-chat", "messages": messages, "stream": False, "max_tokens": 1200}).encode("utf-8") try: conn = (http.client.HTTPSConnection if base.scheme == "https" else http.client.HTTPConnection)(base.hostname, base.port, timeout=timeout) conn.request("POST", (base.path.rstrip("/") or "") + "/chat/completions", body, {"Content-Type": "application/json", "Authorization": "Bearer " + key}) r = conn.getresponse(); raw = r.read(); conn.close() except socket.timeout as e: raise UpstreamError("upstream_timeout", str(e)) except (OSError, http.client.HTTPException) as e: raise UpstreamError("upstream_net", str(e)) if r.status == 401 or r.status == 403: raise UpstreamError("upstream_401", raw[:200].decode("utf-8", "replace")) if r.status == 429: raise UpstreamError("upstream_429", raw[:200].decode("utf-8", "replace")) if r.status >= 500: raise UpstreamError("upstream_5xx", f"{r.status} " + raw[:200].decode("utf-8", "replace")) if r.status != 200: raise UpstreamError("upstream_bad", f"{r.status} " + raw[:200].decode("utf-8", "replace")) try: return json.loads(raw)["choices"][0]["message"]["content"] except Exception as e: raise UpstreamError("upstream_bad", f"{type(e).__name__}: {raw[:200].decode('utf-8', 'replace')}") # ── 核心 (纯函数, transport 可注入) ───────────────────────────────── def ask(q, user, transport=None, ip="-"): """→ {ok, mode, answer|None, err|None, err_code|None, refs, model, secs, request_id}. 任何失败 answer 恒 None; 绝不在失败时给出貌似回答的文本.""" rid = uuid.uuid4().hex[:12]; t0 = time.time(); q = (q or "").strip(); qsha = hashlib.sha256(q.encode("utf-8")).hexdigest()[:16] base = dict(request_id=rid, refs=[], model=os.environ.get("DEEPSEEK_MODEL") or "deepseek-chat", answer=None, err=None, err_code=None) def done(**kw): out = dict(base, **kw); out["secs"] = round(time.time() - t0, 2) audit("qa", request_id=rid, user=user.get("name"), role=user.get("role"), ip=ip, q_sha16=qsha, q_len=len(q), mode=out["mode"], ok=out["ok"], err_code=out["err_code"], secs=out["secs"], n_refs=len(out["refs"])) return out if not q: return done(ok=False, mode=mode_now(), err="问题为空", err_code="empty") if len(q) > LIMITS["max_q_chars"]: return done(ok=False, mode=mode_now(), err=f"问题超长 (>{LIMITS['max_q_chars']} 字)", err_code="too_long") refs = pick_refs(q); base["refs"] = [c["claim_id"] for c in refs] if not api_key(): if os.environ.get("GUANLAN_QA_DEMO") == "1": demo = "[演示回答] 云端模型未配置, 以下不是模型输出, 只是契约里与问题最接近的结论原文:\n" + ("\n".join(f"- {c['claim_id']} [{c['verdict']}] {c['title_public'][:160]}" for c in refs) or "- (契约里没有相关结论)") return done(ok=True, mode="演示回答", answer=demo, model=None) return done(ok=False, mode="不可用", err=ERR_TEXT["no_key"], err_code="no_key", model=None) try: txt = (transport or deepseek_transport)(build_messages(q, refs)) except UpstreamError as e: print(f"[qa] {rid} upstream {e.code}: {e.detail}", file=sys.stderr, flush=True) # 上游原文只进服务端 stderr (已去密钥), 不进回包 return done(ok=False, mode="不可用", err=ERR_TEXT.get(e.code, ERR_TEXT["upstream_bad"]), err_code=e.code) except Exception as e: print(f"[qa] {rid} internal {type(e).__name__}: {redact(str(e))[:200]}", file=sys.stderr, flush=True) return done(ok=False, mode="不可用", err="服务端内部错误, 本次未得到回答", err_code="internal") txt = re.sub(r".*?", "", str(txt or ""), flags=re.S).strip() if not txt: return done(ok=False, mode="不可用", err=ERR_TEXT["upstream_bad"], err_code="upstream_bad") return done(ok=True, mode="云端 DeepSeek", answer=DP.scrub(redact(txt))[:LIMITS["max_answer_chars"]]) # ── HTTP ────────────────────────────────────────────────────────── USERS = None; SESS = Sessions(); RL = RateLimiter(); COOKIE = "guanlan_sid" class H(BaseHTTPRequestHandler): server_version = "guanlan-qa/0.1"; sys_version = "" def log_message(self, *a): pass def _json(self, obj, code=200, extra=None): b = json.dumps(obj, ensure_ascii=False).encode("utf-8") self.send_response(code); self.send_header("Content-Type", "application/json; charset=utf-8"); self.send_header("Content-Length", str(len(b))) self.send_header("Cache-Control", "no-store"); self.send_header("X-Content-Type-Options", "nosniff") for k, v in (extra or {}).items(): self.send_header(k, v) self.end_headers(); self.wfile.write(b) def _ip(self): return self.headers.get("X-Forwarded-For", self.client_address[0]).split(",")[0].strip() def _sid(self): c = http.cookies.SimpleCookie(self.headers.get("Cookie", "")); return c[COOKIE].value if COOKIE in c else None def _user(self): return SESS.get(self._sid()) def _body(self): n = int(self.headers.get("Content-Length") or 0) if n > LIMITS["max_body_bytes"]: return None try: return json.loads(self.rfile.read(n) or b"{}") except Exception: return {} def do_GET(self): p = urllib.parse.urlparse(self.path).path if p == "/healthz": return self._json({"ok": True, "mode": mode_now(), "service": "guanlan-qa"}) u = self._user() if not u: return self._json({"ok": False, "err": "未登录", "err_code": "unauthorized", "mode": mode_now()}, 401) if p == "/api/me": return self._json({"ok": True, "user": u["name"], "role": u["role"], "expires": u["expires"], "mode": mode_now()}) if p == "/api/qa/status": return self._json({"ok": True, "mode": mode_now(), "model": (os.environ.get("DEEPSEEK_MODEL") or "deepseek-chat") if api_key() else None, "limits": {k: v for k, v in LIMITS.items() if k in ("max_q_chars", "rate_user", "upstream_timeout_s")}, "refs_available": bool(claims_public())}) if p == "/api/admin/users": if u["role"] != "admin": audit("forbidden", user=u["name"], path=p, ip=self._ip()); return self._json({"ok": False, "err": "权限不足", "err_code": "forbidden"}, 403) return self._json({"ok": True, "users": USERS.list()}) return self._json({"ok": False, "err": "not found"}, 404) def do_POST(self): p = urllib.parse.urlparse(self.path).path; ip = self._ip() ok, wait = RL.allow("ip:" + ip, *LIMITS["rate_ip"]) if not ok: audit("rate_limited", ip=ip, path=p); return self._json({"ok": False, "err": "请求过于频繁", "err_code": "rate_limited", "mode": mode_now()}, 429, {"Retry-After": str(wait)}) body = self._body() if body is None: return self._json({"ok": False, "err": "请求体过大", "err_code": "too_large", "mode": mode_now()}, 413) if p == "/api/login": u, why = USERS.verify(body.get("name"), body.get("password")); audit("login", user=body.get("name"), ok=why == "ok", reason=why, ip=ip) if not u: return self._json({"ok": False, "err": "用户名或密码错误, 或账号已过期/撤销", "err_code": "unauthorized"}, 401) tok = SESS.create(u); flags = "; HttpOnly; SameSite=Strict; Path=/" + ("" if os.environ.get("GUANLAN_QA_INSECURE_COOKIE") == "1" else "; Secure") return self._json({"ok": True, "user": u["name"], "role": u["role"], "expires": u["expires"], "mode": mode_now()}, 200, {"Set-Cookie": f"{COOKIE}={tok}{flags}"}) u = self._user() if not u: return self._json({"ok": False, "err": "未登录", "err_code": "unauthorized", "mode": mode_now()}, 401) if p == "/api/logout": SESS.drop(self._sid()); audit("logout", user=u["name"], ip=ip); return self._json({"ok": True}, 200, {"Set-Cookie": f"{COOKIE}=; Max-Age=0; Path=/"}) if p == "/api/qa": ok, wait = RL.allow("user:" + u["name"], *LIMITS["rate_user"]) if not ok: audit("rate_limited", user=u["name"], ip=ip); return self._json({"ok": False, "err": "本账号提问过于频繁, 请稍后", "err_code": "rate_limited", "mode": mode_now()}, 429, {"Retry-After": str(wait)}) r = ask(body.get("q"), u, ip=ip) return self._json(r, 200 if r["ok"] else {"too_long": 413, "empty": 400}.get(r["err_code"], 503)) return self._json({"ok": False, "err": "not found"}, 404) def selftest(): """不联网: 用桩走一遍 无密钥 / 演示 / 401 / 429 / 超时 / 5xx / 成功, 打印每种模式标签.""" saved = {k: os.environ.get(k) for k in ("DEEPSEEK_API_KEY", "GUANLAN_QA_DEMO")}; u = {"name": "selftest", "role": "user"}; out = [] try: os.environ.pop("DEEPSEEK_API_KEY", None); os.environ.pop("GUANLAN_QA_DEMO", None); out.append(("无密钥", ask("叶根螺栓", u))) os.environ["GUANLAN_QA_DEMO"] = "1"; out.append(("演示", ask("叶根螺栓", u))); os.environ.pop("GUANLAN_QA_DEMO") os.environ["DEEPSEEK_API_KEY"] = "sk-selftest-not-a-real-key-000000" for code in ("upstream_401", "upstream_429", "upstream_timeout", "upstream_5xx"): def bad(m, code=code): raise UpstreamError(code, "Bearer sk-selftest-not-a-real-key-000000 detail") out.append((code, ask("叶根螺栓", u, transport=bad))) out.append(("成功", ask("叶根螺栓", u, transport=lambda m: "如东 WTG31 叶根螺栓断裂仍在发生 (RD-2026-08-003)"))) finally: for k, v in saved.items(): if v is None: os.environ.pop(k, None) else: os.environ[k] = v for n, r in out: print(f"{n:16} mode={r['mode']:12} ok={r['ok']!s:5} err_code={r['err_code']} answer={(r['answer'] or '')[:60]!r} refs={r['refs'][:3]}") return out def main(): global USERS ap = argparse.ArgumentParser(); sub = ap.add_subparsers(dest="cmd", required=True) s = sub.add_parser("serve"); s.add_argument("--port", type=int, default=8090); s.add_argument("--bind", default="127.0.0.1") uu = sub.add_parser("user"); uu.add_argument("op", choices=["add", "revoke", "list"]); uu.add_argument("name", nargs="?"); uu.add_argument("--role", default="user", choices=["user", "admin"]); uu.add_argument("--days", type=int, default=LIMITS["temp_pw_days"]) sub.add_parser("selftest"); a = ap.parse_args(); USERS = Users() if a.cmd == "user": if a.op == "add": print(f"账号 {a.name} ({a.role}) 临时密码 (只显示这一次, {a.days} 天有效):", USERS.add(a.name, a.role, a.days)); return 0 if a.op == "revoke": USERS.revoke(a.name); print("已撤销", a.name); return 0 for r in USERS.list(): print(json.dumps(r, ensure_ascii=False)) return 0 if a.cmd == "selftest": selftest(); return 0 print(f"[qa] :{a.port} mode={mode_now()} users={len(USERS.list())} refs={len(claims_public())} audit={AUDIT_PATH}", file=sys.stderr, flush=True) if not api_key(): print("[qa] 警告: 未配置 DEEPSEEK_API_KEY, /api/qa 全部返回 不可用" + (" (演示回答)" if mode_now() == "演示回答" else ""), file=sys.stderr, flush=True) ThreadingHTTPServer((a.bind, a.port), H).serve_forever() if __name__ == "__main__": sys.exit(main() or 0)