feat(proxy): T-X11 采纳 cortiq shadow 旁路评审——live T1 本地应答异步云端评审喂晋升表

- routes:live 分流 T1 响应后按 sample_rate 采样旁路——premium 条目以评审提示词
  给本地答案打 PASS/FAIL,经 PromotionTable.observe 写入 T-X5 晋升表
  (打通晋升质量信号的自动来源,夜间 labeler 之外的实时通路)
- 零客户端延迟:asyncio.create_task 旁路 + wait_for 超时 + Semaphore(2) 并发上限;
  上游故障/解析失败记 ERROR,不算 FAIL(评审器不可用不惩罚本地模型);
  仅非流式响应采样(流式无同步答案文本),默认关闭(sense.shadow_review.enabled)
- 模块级计数字典 sampled/pass/fail/error + reset_shadow_review_stats() 测试钩子
- 新增 tests/test_shadow_review.py 5 项(默认关/条目选取/PASS-FAIL-ERROR 解析/
  晋升表喂入/关闭态零副作用),全假上游不依赖真实模型
This commit is contained in:
tzt
2026-09-19 10:03:38 +08:00
parent e5470a3715
commit f04a6c4c43
2 changed files with 284 additions and 0 deletions
+156
View File
@@ -117,6 +117,10 @@ def build_proxy_router(cfg: ProxyConfig, pool, settings_provider=None,
request_id=request_id)
if tier_used:
resp.headers["x-campus-tier"] = tier_used
# T-X11(采纳 cortiq shadow 旁路评审):live T1 本地应答按采样率
# 异步旁路云端评审,喂晋升表(零客户端延迟,任何异常不影响响应)
if tier_used == "T1":
_maybe_shadow_review(body, resp, pool, settings_provider, scfg)
# T1 审计抽样(§9.4):按 review.sample_rate 入队人工核;
# promo 标签供人工 verdict 回填晋升表(T-X5)。
# 2026-09 修复:原块引用未定义的 load_config/get_review/request_id
@@ -454,6 +458,158 @@ def reset_sense_runtimes() -> None:
_sense_runtimes = {}
# ---------------------------------------------------------------------------
# T-X11(采纳 cortiq shadow 旁路评审设计):live T1 本地应答的异步云端评审。
# 按采样率旁路:premium 条目给本地答案打 PASS/FAIL,写入 T-X5 晋升表——
# 打通晋升表的质量信号自动来源。零客户端延迟(asyncio.create_task 旁路)。
# ---------------------------------------------------------------------------
_shadow_review_stats: Dict[str, int] = {"sampled": 0, "pass": 0, "fail": 0,
"error": 0, "skipped_busy": 0}
_shadow_sem: Optional[asyncio.Semaphore] = None
_SHADOW_CONCURRENCY = 2 # 旁路评审并发上限(保护上游)
_SHADOW_MAX_ANSWER_CHARS = 4000 # 送审答案截断
def reset_shadow_review_stats() -> None:
"""测试用:清零旁路评审计数并复位并发闸门。"""
global _shadow_review_stats, _shadow_sem
_shadow_review_stats = {"sampled": 0, "pass": 0, "fail": 0,
"error": 0, "skipped_busy": 0}
_shadow_sem = None
def _shadow_review_settings(settings_dict: dict) -> Dict[str, Any]:
"""解析 sense.shadow_review 设置(默认关闭;fail-safe 兜底值)。"""
raw = (settings_dict.get("sense") or {}).get("shadow_review") or {}
if not isinstance(raw, dict):
raw = {}
try:
rate = max(0.0, min(1.0, float(raw.get("sample_rate", 0.1) or 0.1)))
except (TypeError, ValueError):
rate = 0.1
try:
timeout_s = max(1.0, float(raw.get("timeout_s", 20) or 20))
except (TypeError, ValueError):
timeout_s = 20.0
return {"enabled": bool(raw.get("enabled", False)),
"sample_rate": rate, "timeout_s": timeout_s,
"tier": str(raw.get("tier", "premium") or "premium")}
def _pick_review_entry(pool, tier_hint: str) -> Optional[Dict[str, Any]]:
"""选评审条目:优先指定档位,回退任意启用真实后端条目。"""
entries = [e for e in _pool_entries(pool)
if e.get("enabled") and e.get("backend") not in ("mock",)
and e.get("base_url")]
for e in entries:
if (e.get("tier") or "") == tier_hint:
return e
return entries[0] if entries else None
async def _grade_answer(entry, question: str, answer: str, model: str,
timeout_s: float) -> str:
"""用指定池条目评审本地答案,返回 PASS / FAIL / ERROR(上游故障不抛出)。"""
judge_body = {
"model": model,
"messages": [
{"role": "system", "content": "你是严格的答案质量评审员。只输出 PASS 或 FAIL。"},
{"role": "user",
"content": (f"问题:{question}\n\n候选答案:{answer[:_SHADOW_MAX_ANSWER_CHARS]}\n\n"
"该答案是否正确且切题?只输出一个词:PASS 或 FAIL。")},
],
"max_tokens": 8,
}
parts: List[str] = []
sink: Dict[str, Any] = {}
async def _collect():
async for raw in upstream_stream(judge_body, entry, sink, [entry]):
line = raw.decode("utf-8", errors="replace").strip()
for sub in line.split("\n\n"):
if not sub.startswith("data:"):
continue
payload = sub[5:].strip()
if not payload or payload == "[DONE]":
continue
try:
obj = json.loads(payload)
except json.JSONDecodeError:
continue
delta = (obj.get("choices") or [{}])[0].get("delta") or {}
if delta.get("content"):
parts.append(str(delta["content"]))
try:
await asyncio.wait_for(_collect(), timeout=timeout_s)
except Exception: # noqa: BLE001 超时/上游异常:评审失败不算 FAIL
return "ERROR"
text = "".join(parts).strip().upper()
if "PASS" in text:
return "PASS"
if "FAIL" in text:
return "FAIL"
return "ERROR"
async def _shadow_review_task(question: str, answer: str, pool,
settings_dict: dict, promotion,
tier_used: str = "T1") -> None:
"""旁路评审任务体:调用方 create_task 丢弃即可,任何异常不影响主链路。"""
cfgs = _shadow_review_settings(settings_dict)
entry = _pick_review_entry(pool, cfgs["tier"])
if entry is None:
return
global _shadow_sem
if _shadow_sem is None:
_shadow_sem = asyncio.Semaphore(_SHADOW_CONCURRENCY)
try:
async with _shadow_sem:
verdict = await _grade_answer(entry, question, answer,
str(entry.get("model") or ""),
cfgs["timeout_s"])
except Exception: # noqa: BLE001
verdict = "ERROR"
_shadow_review_stats["sampled"] += 1
if verdict in ("PASS", "FAIL"):
_shadow_review_stats["pass" if verdict == "PASS" else "fail"] += 1
try:
from gateway.sense.promotion import promotion_label
promotion.observe(promotion_label("proxy", ""),
ok=(verdict == "PASS"), tier=tier_used)
except Exception: # noqa: BLE001
pass
else:
_shadow_review_stats["error"] += 1
def _maybe_shadow_review(body: dict, resp, pool, settings_provider, scfg) -> None:
"""主链路挂钩:live T1 且命中采样率时旁路评审(仅非流式响应可取答案文本)。"""
try:
import random as _random
cfgs = _shadow_review_settings(settings_provider().to_dict())
if not cfgs["enabled"] or _random.random() >= cfgs["sample_rate"]:
return
if isinstance(resp, StreamingResponse):
return # 流式响应此刻无完整答案文本,跳过(后续迭代可经 sink 采集)
payload = json.loads(bytes(resp.body))
answer = str((((payload.get("choices") or [{}])[0])
.get("message") or {}).get("content") or "")
if not answer:
return
question = "\n".join(str(m.get("content") or "")
for m in (body.get("messages") or []))
_, grader = _get_sense_runtime(scfg)
promotion = grader._promotion # 同包内复用晋升表(T-X5
if promotion is None:
return
asyncio.get_running_loop().create_task(
_shadow_review_task(question, answer, pool,
settings_provider().to_dict(), promotion))
except Exception: # noqa: BLE001 旁路失败不影响主响应(D-G4 纪律)
pass
def _route_sig(needs: Dict[str, Any], model: str, cfg: ProxyConfig) -> str:
"""路由签名(T-X10,采纳 cortiq 语义缓存路由签名分桶)。
+128
View File
@@ -0,0 +1,128 @@
"""shadow 旁路评审单元测试(T-X11,采纳 cortiq shadow bypass 设计)。
全部走 monkeypatch 假上游,不依赖真实模型/API key。
"""
import asyncio
import json
import pytest
import gateway.proxy.routes as R
from gateway.proxy.routes import (_grade_answer, _pick_review_entry,
_shadow_review_settings,
reset_shadow_review_stats)
from gateway.sense.promotion import PromotionTable, promotion_label
@pytest.fixture(autouse=True)
def _reset():
reset_shadow_review_stats()
yield
reset_shadow_review_stats()
def _fake_upstream(responses):
"""构造假 upstream_stream:按调用顺序吐预定 verdict 文本(SSE 形态)。"""
calls = {"n": 0}
def _stream(body, entry, sink, chain):
async def gen():
idx = min(calls["n"], len(responses) - 1)
calls["n"] += 1
verdict = responses[idx]
chunk = {"choices": [{"delta": {"content": verdict}}]}
yield (f"data: {json.dumps(chunk)}\n\n").encode("utf-8")
yield b"data: [DONE]\n\n"
return gen()
return _stream, calls
def _entry(tier="premium", model="cloud-x"):
return {"id": "e1", "enabled": True, "backend": "openai",
"base_url": "http://127.0.0.1:9/v1", "model": model, "tier": tier}
def test_settings_default_disabled_and_clamped():
"""默认关闭;采样率/超时非法值被夹取。"""
s = _shadow_review_settings({})
assert s["enabled"] is False and s["sample_rate"] == 0.1
assert _shadow_review_settings({"sense": {"shadow_review": {
"enabled": True, "sample_rate": 5, "timeout_s": -3}}})["sample_rate"] == 1.0
assert _shadow_review_settings({"sense": {"shadow_review": {
"enabled": True, "timeout_s": -3}}})["timeout_s"] == 1.0
def test_pick_review_entry_prefers_tier():
"""评审条目优先指定档位,缺失时回退任意启用条目。"""
pool_entries = [_entry("budget", "b1"), _entry("premium", "p1")]
class _P:
@staticmethod
def usable_entries():
return list(pool_entries)
assert _pick_review_entry(_P, "premium")["model"] == "p1"
assert _pick_review_entry(_P, "local")["model"] in ("b1", "p1")
def test_grade_answer_parses_pass_fail_and_error(monkeypatch):
"""PASS/FAIL 解析;上游异常回 ERROR 而非抛出。"""
fake, calls = _fake_upstream(["PASS", "FAIL"])
monkeypatch.setattr(R, "upstream_stream", fake)
e = _entry()
assert asyncio.run(_grade_answer(e, "", "", "cloud-x", 5)) == "PASS"
assert asyncio.run(_grade_answer(e, "", "", "cloud-x", 5)) == "FAIL"
def _boom(*a, **k):
async def gen():
raise RuntimeError("上游炸了")
yield b"" # pragma: no cover
return gen()
monkeypatch.setattr(R, "upstream_stream", _boom)
assert asyncio.run(_grade_answer(e, "", "", "cloud-x", 5)) == "ERROR"
assert calls["n"] == 2
def test_shadow_review_task_feeds_promotion(tmp_path, monkeypatch):
"""PASS/FAIL 写入晋升表(observe 计数),ERROR 只计错误不进通过率。"""
from gateway.sense.store import SenseStore
sstore = SenseStore.init_db(str(tmp_path / "sense.sqlite3"))
promo = PromotionTable(sstore, now=lambda: 1_000_000.0)
label = promotion_label("proxy", "")
fake, _ = _fake_upstream(["PASS", "PASS", "FAIL", "垃圾输出"])
monkeypatch.setattr(R, "upstream_stream", fake)
entry = _entry()
sd = {"sense": {"shadow_review": {"enabled": True, "sample_rate": 1.0}}}
async def _run():
for _ in range(4):
await R._shadow_review_task("问题", "本地答案", _FixedPool(entry),
sd, promo)
asyncio.run(_run())
assert R._shadow_review_stats == {"sampled": 4, "pass": 2, "fail": 1,
"error": 1, "skipped_busy": 0}
row = sstore.get_promotion(label)
assert row["n_total"] == 3 and row["n_ok"] == 2 # ERROR 不进通过率
class _FixedPool:
def __init__(self, entries):
self._e = entries if isinstance(entries, list) else [entries]
def usable_entries(self):
return list(self._e)
def test_maybe_shadow_review_disabled_is_noop(monkeypatch):
"""默认关闭:不建任务、不计采样。"""
class _SP:
@staticmethod
def to_dict():
return {"sense": {"shadow_review": {"enabled": False}}}
payload = json.dumps({"choices": [{"message": {"content": "本地答案"}}]}).encode()
resp = type("JR", (), {"body": payload})()
before = dict(R._shadow_review_stats)
R._maybe_shadow_review({}, resp, _FixedPool(_entry()), _SP, scfg=None)
assert R._shadow_review_stats == before