diff --git a/gateway/proxy/routes.py b/gateway/proxy/routes.py index 932cd9f..00af69a 100644 --- a/gateway/proxy/routes.py +++ b/gateway/proxy/routes.py @@ -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 语义缓存路由签名分桶)。 diff --git a/tests/test_shadow_review.py b/tests/test_shadow_review.py new file mode 100644 index 0000000..444dc7b --- /dev/null +++ b/tests/test_shadow_review.py @@ -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