From 1aaa05cacc53a3e64f2e3aa2e723dc04f917cf2d Mon Sep 17 00:00:00 2001 From: tzt <14718231+flying-travel@user.noreply.gitee.com> Date: Fri, 18 Sep 2026 23:01:52 +0800 Subject: [PATCH] =?UTF-8?q?feat(sense):=20T-X5=20=E4=BA=91=E7=AB=AF?= =?UTF-8?q?=E8=AF=84=E5=88=A4=E6=99=8B=E5=8D=87=E8=A1=A8=C3=97conformal=20?= =?UTF-8?q?=E5=8F=8C=E9=97=B8=E9=97=A8=E2=80=94=E2=80=94=E6=9C=AC=E5=9C=B0?= =?UTF-8?q?=E6=A1=A3=E8=87=AA=E5=8A=A8=E6=8E=A5=E7=AE=A1=EF=BC=88=E9=87=87?= =?UTF-8?q?=E7=BA=B3=20cortiq=20promotion=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - gateway/sense/promotion.py:PromotionTable 状态机(candidate/promoted/demoted) * 晋升门:n_total>=n_min 且 通过率>=promote_lb 且 soak 浸泡期满足 * 退化:promoted 期间通过率 标签 -> 人工 verdict approve=ok 经 /review/{id} 提交时回填晋升表(失败不影响审核主流程) - 双闸门语义:conformal 阈值保证单条决策风险率,晋升表保证标签级接管节奏; collect 数据不足(n_min 未满)不晋升,全部迁移可经 T-X4 留痕审计 pytest 474 passed(T-X4 后 466 + 8) --- gateway/api.py | 16 +++ gateway/proxy/routes.py | 11 ++- gateway/sense/grader.py | 19 +++- gateway/sense/promotion.py | 81 ++++++++++++++++ gateway/sense/store.py | 32 ++++++ tests/test_sense_promotion.py | 178 ++++++++++++++++++++++++++++++++++ 任务拆解与执行计划.md | 2 + 7 files changed, 334 insertions(+), 5 deletions(-) create mode 100644 gateway/sense/promotion.py create mode 100644 tests/test_sense_promotion.py diff --git a/gateway/api.py b/gateway/api.py index cf07e6b..4e76b36 100644 --- a/gateway/api.py +++ b/gateway/api.py @@ -451,6 +451,22 @@ try: raise HTTPException(status_code=400, detail=str(e)) if not ok: raise HTTPException(status_code=404, detail=f"审核记录不存在或已审核: {review_id}") + # 晋升表回填(T-X5):带 promo 标签的 sense 审计样本,人工 verdict + # approve=ok -> 晋升表累计;失败不影响审核主流程 + try: + row = get_review().get(review_id) + promo = next((t[len("promo:"):] for t in (row.get("tags") or []) + if isinstance(t, str) and t.startswith("promo:")), None) + if promo: + from gateway.sense.config import build_sense_config + from gateway.sense.promotion import PromotionTable + from gateway.sense.store import SenseStore + scfg = build_sense_config(settings_store().to_dict()) + if scfg.enabled: + PromotionTable(SenseStore.init_db(scfg.db_path)).observe( + promo, ok=(verdict == "approve")) + except Exception: + pass return {"ok": True, "review_id": review_id, "verdict": verdict} # ---------------- 模型池(多价位异构模型) ---------------- diff --git a/gateway/proxy/routes.py b/gateway/proxy/routes.py index 15e8b94..d7c7c62 100644 --- a/gateway/proxy/routes.py +++ b/gateway/proxy/routes.py @@ -93,11 +93,13 @@ def build_proxy_router(cfg: ProxyConfig, pool, settings_provider=None) -> APIRou from gateway.sense.config import build_sense_config from gateway.sense.grader import Grader from gateway.sense.observer import get_observer + from gateway.sense.promotion import PromotionTable from gateway.sense.store import SenseStore scfg = build_sense_config(settings_provider().to_dict()) if scfg.enabled and scfg.mode == "live": sstore = SenseStore.init_db(scfg.db_path) - grader = Grader(scfg, sstore, get_observer(sstore)) + grader = Grader(scfg, sstore, get_observer(sstore), + promotion=PromotionTable(sstore)) qtext = "\n".join(str(m.get("content") or "") for m in (body.get("messages") or [])) d = await grader.decide(qtext or str(body.get("model") or ""), @@ -114,14 +116,17 @@ def build_proxy_router(cfg: ProxyConfig, pool, settings_provider=None) -> APIRou resp = await _run_chat(body, dict(request.headers), ctx, cfg, ledger, pool) if tier_used: resp.headers["x-campus-tier"] = tier_used - # T1 审计抽样(§9.4):按 review.sample_rate 入队人工核 + # T1 审计抽样(§9.4):按 review.sample_rate 入队人工核; + # promo 标签供人工 verdict 回填晋升表(T-X5) try: import random as _random + from gateway.sense.promotion import promotion_label rate = float(load_config().get("review", {}).get("sample_rate", 0.1)) if _random.random() < rate: get_review().enqueue( request_id + "-sense", str(body.get("model") or "proxy"), - "(sense T1 审计抽样)", tags=["sense_t1"], + "(sense T1 审计抽样)", + tags=["sense_t1", "promo:" + promotion_label("proxy", "")], reason="sense_audit") except Exception: pass diff --git a/gateway/sense/grader.py b/gateway/sense/grader.py index 724d514..00de8ba 100644 --- a/gateway/sense/grader.py +++ b/gateway/sense/grader.py @@ -21,6 +21,7 @@ from gateway.sense.decision_cache import DecisionCache from gateway.sense.embedder import embed from gateway.sense.errors import EmbedderDown from gateway.sense.features import gate +from gateway.sense.promotion import promotion_label from gateway.sense.store import json_dumps @@ -99,10 +100,11 @@ def _rejected_tiers(fallback: bool, probs: Dict[str, float], class Grader: """决策器(持有 active 工件缓存;工件/阈值切换后调 invalidate)。""" - def __init__(self, cfg: SenseConfig, store, observer=None): + def __init__(self, cfg: SenseConfig, store, observer=None, promotion=None): self.cfg = cfg self.store = store self.observer = observer + self._promotion = promotion # T-X5:晋升表(可 None) self._head: Optional[LinearHead] = None self._head_loaded = False self._thresholds: Optional[Dict[str, Any]] = None @@ -194,6 +196,18 @@ class Grader: else: tier = "T2" + # ---- 晋升表应用(T-X5:conformal 之外的统计闸门;仅 live)---- + # 仅对 T2 决策且无任何 T3 信号时,允许已晋升标签升级 T1; + # T1/T3 决策与 fallback 路径不受晋升表影响。 + promotion_applied = False + if (self._promotion is not None and self.cfg.mode == "live" + and not fallback and tier == "T2" + and not feats.repo_signals + and probs.get("t3", 0.0) < float(th.get("t3", 0.9))): + if self._promotion.is_promoted(promotion_label(consumer, domain)): + tier = "T1" + promotion_applied = True + # ---- mode 裁剪(D-G7)---- if self.cfg.mode == "live": executed = tier # live:决策即执行 @@ -208,7 +222,8 @@ class Grader: "intent_blocked": feats.intent_blocked, "repo_signals": feats.repo_signals, "over_length": feats.over_length, - "t1_hard_ok": feats.t1_hard_ok}, + "t1_hard_ok": feats.t1_hard_ok, + "promotion_applied": promotion_applied}, thresholds_version=th_version, head_version=head_version, mode=self.cfg.mode, fallback=fallback, executed_tier=executed) diff --git a/gateway/sense/promotion.py b/gateway/sense/promotion.py new file mode 100644 index 0000000..22e868d --- /dev/null +++ b/gateway/sense/promotion.py @@ -0,0 +1,81 @@ +"""云端评判式晋升表(T-X5,采纳 cortiq promotion.rs,与 conformal 双闸门)。 + +定位:conformal 阈值给出「分对率 >= 1-α」的统计保证(单条查询粒度), +本表在其之上给出「某类任务(label)可由本地档接管」的**节奏自动化**: + +- 观察:人工核验 verdict(approve=ok)等质量信号按 label 累计通过率; +- 晋升:n_total >= n_min 且 通过率 >= promote_lb 且 soak 浸泡期满足 + -> state=promoted(此后 grader 对该 label 的 T2 决策可升级 T1); +- 退化:promoted 期间通过率 < demote_lb -> demoted(自动回退分级路径); +- 一票否决:tier=T3(高复杂度/仓库级)观察不进通过率统计,且若已晋升 + 立即降级——与 cortiq "HIGH tier escalation is served by the cloud" 同构。 + +灰度纪律(沿用 D-G7):collect 数据不足(n_min 未满)不得晋升; +全部状态迁移可由 T-X4 的决策留痕审计。 +""" +from __future__ import annotations + +import time +from typing import Any, Dict, Optional + + +def promotion_label(consumer: str, domain: str = "") -> str: + """晋升粒度:消费方×领域(domain 缺省归 general;调用方可传更细粒度)。""" + return f"{consumer or 'default'}:{domain or 'general'}" + + +class PromotionTable: + """标签晋升状态机(candidate / promoted / demoted;持久化于 sense_promotion)。""" + + def __init__(self, store, now=None, n_min: int = 20, promote_lb: float = 0.95, + soak_days: float = 3.0, demote_lb: float = 0.85, + high_tier_veto: bool = True): + self.store = store + self._now = now or (lambda: time.time()) + self.n_min = max(1, int(n_min)) + self.promote_lb = float(promote_lb) + self.demote_lb = float(demote_lb) + self.soak_s = float(soak_days) * 86400.0 + self.high_tier_veto = bool(high_tier_veto) + + def observe(self, label: str, ok: bool, tier: str = "", + ts: Optional[float] = None) -> Dict[str, Any]: + """记录一次质量信号并推进状态机;返回迁移后的行。""" + ts = float(ts) if ts is not None else float(self._now()) + row = self.store.get_promotion(label) or { + "label": label, "n_total": 0, "n_ok": 0, "state": "candidate", + "first_ts": int(ts), "promoted_ts": 0, "last_ts": 0} + + # 一票否决:T3 观察不进通过率统计;已晋升者立即降级 + if self.high_tier_veto and str(tier).upper() == "T3": + row["last_ts"] = int(ts) + if row["state"] == "promoted": + row["state"] = "demoted" + row["promoted_ts"] = 0 + self.store.upsert_promotion(row) + return row + + row["n_total"] = int(row["n_total"]) + 1 + row["n_ok"] = int(row["n_ok"]) + (1 if ok else 0) + row["last_ts"] = int(ts) + rate = row["n_ok"] / row["n_total"] if row["n_total"] else 0.0 + soaked = (ts - int(row["first_ts"])) >= self.soak_s + + if row["state"] == "promoted": + if row["n_total"] >= self.n_min and rate < self.demote_lb: + row["state"] = "demoted" # 退化自动降级 + row["promoted_ts"] = 0 + else: + # candidate / demoted 共用同一晋升门(降级后可凭数据恢复) + if row["n_total"] >= self.n_min and rate >= self.promote_lb and soaked: + row["state"] = "promoted" + row["promoted_ts"] = int(ts) + self.store.upsert_promotion(row) + return row + + def is_promoted(self, label: str) -> bool: + row = self.store.get_promotion(label) + return bool(row) and row["state"] == "promoted" + + def snapshot(self, label: str) -> Optional[Dict[str, Any]]: + return self.store.get_promotion(label) diff --git a/gateway/sense/store.py b/gateway/sense/store.py index 89dd729..28be37c 100644 --- a/gateway/sense/store.py +++ b/gateway/sense/store.py @@ -41,6 +41,14 @@ CREATE TABLE IF NOT EXISTS sense_decisions( candidate_scores TEXT NOT NULL DEFAULT '{}', rejected TEXT NOT NULL DEFAULT '[]'); CREATE INDEX IF NOT EXISTS idx_dec_ts ON sense_decisions(ts); + +CREATE TABLE IF NOT EXISTS sense_promotion( + label TEXT PRIMARY KEY, n_total INTEGER NOT NULL DEFAULT 0, + n_ok INTEGER NOT NULL DEFAULT 0, + state TEXT NOT NULL DEFAULT 'candidate', + first_ts INTEGER NOT NULL DEFAULT 0, + promoted_ts INTEGER NOT NULL DEFAULT 0, + last_ts INTEGER NOT NULL DEFAULT 0); """ @@ -178,6 +186,30 @@ class SenseStore: out.append(d) return out + # ---------- 晋升表(T-X5:本地档自动接管的双闸门之一) ---------- + def upsert_promotion(self, row: Dict[str, Any]) -> None: + with self._lock, self._connect() as conn: + conn.execute( + """INSERT OR REPLACE INTO sense_promotion + (label, n_total, n_ok, state, first_ts, promoted_ts, last_ts) + VALUES (?,?,?,?,?,?,?)""", + (row["label"], int(row.get("n_total", 0)), + int(row.get("n_ok", 0)), str(row.get("state", "candidate")), + int(row.get("first_ts", 0)), int(row.get("promoted_ts", 0)), + int(row.get("last_ts", 0)))) + + def get_promotion(self, label: str) -> Optional[Dict[str, Any]]: + with self._lock, self._connect() as conn: + row = conn.execute( + "SELECT * FROM sense_promotion WHERE label = ?", (label,)).fetchone() + return dict(row) if row else None + + def list_promotions(self) -> List[Dict[str, Any]]: + with self._lock, self._connect() as conn: + rows = conn.execute( + "SELECT * FROM sense_promotion ORDER BY last_ts DESC").fetchall() + return [dict(r) for r in rows] + # ---------- 工件 ---------- def register_artifact(self, version: str, kind: str, path: str, metrics: Dict[str, Any], active: bool = False, diff --git a/tests/test_sense_promotion.py b/tests/test_sense_promotion.py new file mode 100644 index 0000000..8249290 --- /dev/null +++ b/tests/test_sense_promotion.py @@ -0,0 +1,178 @@ +"""晋升表测试(T-X5):candidate→promoted→demoted→恢复状态机 + 一票否决 + +grader 晋升应用(T2→T1)与 T3/闸门互斥。""" +import asyncio + +import pytest + +from gateway.sense.config import build_sense_config +from gateway.sense.grader import Grader +from gateway.sense.promotion import PromotionTable, promotion_label +from gateway.sense.store import SenseStore + + +class _Clock: + def __init__(self, t=1789874000.0): + self.t = t + + def __call__(self): + return self.t + + +def _table(tmp_path, **kw): + store = SenseStore.init_db(tmp_path / "s.sqlite3") + clock = _Clock() + table = PromotionTable(store, now=clock, **kw) + return store, clock, table + + +def test_promotion_requires_min_and_soak(tmp_path): + """n_min 未满不晋升; soak 期未满不晋升;两者满足才 promoted。""" + store, clock, table = _table(tmp_path, n_min=5, promote_lb=0.9, + soak_days=3.0) + label = "proxy:general" + for _ in range(5): + table.observe(label, ok=True, ts=clock.t) # 通过率 1.0 但浸泡期未满 + assert not table.is_promoted(label) + clock.t += 4 * 86400.0 # 浸泡期满足 + table.observe(label, ok=True, ts=clock.t) + assert table.is_promoted(label) + + +def test_promotion_demotes_on_degradation(tmp_path): + """promoted 后通过率跌破 demote_lb -> 自动降级。""" + store, clock, table = _table(tmp_path, n_min=5, promote_lb=0.9, + demote_lb=0.85, soak_days=0.0) + label = "proxy:general" + for _ in range(6): + table.observe(label, ok=True, ts=clock.t) + assert table.is_promoted(label) + for _ in range(6): + clock.t += 60 + table.observe(label, ok=False, ts=clock.t) # 通过率跌至 6/12 = 0.5 + assert not table.is_promoted(label) + assert table.snapshot(label)["state"] == "demoted" + + +def test_promotion_recovery_after_demotion(tmp_path): + """demoted 后凭数据恢复:再次满足晋升门 -> promoted。""" + store, clock, table = _table(tmp_path, n_min=4, promote_lb=0.9, + demote_lb=0.85, soak_days=0.0) + label = "proxy:general" + for _ in range(4): + table.observe(label, ok=True, ts=clock.t) + assert table.is_promoted(label) + for _ in range(4): + clock.t += 60 + table.observe(label, ok=False, ts=clock.t) # 跌破 + assert table.snapshot(label)["state"] == "demoted" + for _ in range(20): + clock.t += 60 + table.observe(label, ok=True, ts=clock.t) # 20/24 ≈ 0.83... 继续 + # 24 条中 20 ok = 0.833 < 0.9,再补 ok 至 >0.9 + for _ in range(6): + clock.t += 60 + table.observe(label, ok=True, ts=clock.t) + # 30 条 26 ok ≈ 0.867,仍 <0.9;再补 + for _ in range(6): + clock.t += 60 + table.observe(label, ok=True, ts=clock.t) + # 36 条 32 ok ≈ 0.889;最后一条后 >= 0.9 + table.observe(label, ok=True, ts=clock.t + 60) # 37 条 33 ok ≈ 0.892 + for _ in range(4): + table.observe(label, ok=True, ts=clock.t + 120) + # 41 条 37 ok ≈ 0.902 >= 0.9 + assert table.snapshot(label)["state"] == "promoted" + + +def test_high_tier_veto(tmp_path): + """T3 观察:不计入通过率;已晋升立即降级。""" + store, clock, table = _table(tmp_path, n_min=3, promote_lb=0.9, soak_days=0.0) + label = "proxy:general" + for _ in range(3): + table.observe(label, ok=True, ts=clock.t) + assert table.is_promoted(label) + before = table.snapshot(label)["n_total"] + table.observe(label, ok=True, tier="T3", ts=clock.t + 10) + snap = table.snapshot(label) + assert snap["n_total"] == before # T3 不进统计 + assert snap["state"] == "demoted" # 一票否决 + + +def test_promotion_label_format(): + assert promotion_label("proxy", "") == "proxy:general" + assert promotion_label("proxy", "code") == "proxy:code" + + +# ---------------- grader 晋升应用 ---------------- + +class _FakeHeadMid: + version = "test-head" + + def predict(self, vec): + return {"t1": 0.50, "t2": 0.47, "t3": 0.03} # 低于阈值 -> default T2 + + +class _NoopObserver: + def log(self, obs): + pass + + +def _make(tmp_path, table): + cfg = build_sense_config({"sense": { + "enabled": True, "mode": "live", + "db_path": str(tmp_path / "s.sqlite3"), + "models_dir": str(tmp_path / "models")}}) + store = SenseStore.init_db(cfg.db_path) + g = Grader(cfg, store, _NoopObserver(), promotion=table) + g._load_head = lambda: _FakeHeadMid() + + async def fake_embed(text, c): + return [1, 2, 3] + + import gateway.sense.grader as gr + orig = gr.embed + gr.embed = fake_embed + return g, orig + + +def test_grader_applies_promotion_t2_to_t1(tmp_path): + """已晋升标签:T2 决策升级 T1 且 hard_gates 标注 promotion_applied。""" + store = SenseStore.init_db(tmp_path / "s.sqlite3") + clock = _Clock() + table = PromotionTable(store, now=clock, n_min=2, promote_lb=0.9, soak_days=0.0) + label = promotion_label("proxy", "general") + table.observe(label, ok=True, ts=clock.t) + table.observe(label, ok=True, ts=clock.t) + g, orig = _make(tmp_path, table) + import gateway.sense.grader as gr + try: + d = asyncio.run(g.decide("什么是递归", "proxy")) + finally: + gr.embed = orig + assert d.tier == "T1" + assert d.hard_gates["promotion_applied"] is True + + +def test_grader_skips_promotion_when_not_promoted(tmp_path): + """未晋升:同样的 T2 决策保持 T2。""" + store = SenseStore.init_db(tmp_path / "s2.sqlite3") + table = PromotionTable(store, now=_Clock(), n_min=100, soak_days=0.0) + g, orig = _make(tmp_path, table) + import gateway.sense.grader as gr + try: + d = asyncio.run(g.decide("什么是递归", "proxy")) + finally: + gr.embed = orig + assert d.tier == "T2" + assert d.hard_gates["promotion_applied"] is False + + +def test_grader_promotion_not_applied_without_table(tmp_path): + """未接晋升表(向后兼容):行为与既有决策一致。""" + g, orig = _make(tmp_path, None) + import gateway.sense.grader as gr + try: + d = asyncio.run(g.decide("什么是递归", "proxy")) + finally: + gr.embed = orig + assert d.tier == "T2" diff --git a/任务拆解与执行计划.md b/任务拆解与执行计划.md index a88f504..86304cf 100644 --- a/任务拆解与执行计划.md +++ b/任务拆解与执行计划.md @@ -170,3 +170,5 @@ P0 完成后的能力:干净的后端抽象 + 可量化的评测 + 可追溯 | T-X3 | 路由决策缓存(外部采纳 cortiq 决策哈希缓存):DecisionCache(sha256+LRU+TTL60s,4096 条)前置 grader live 模式,embed+线性头去重;观察/审计不跳过、invalidate 同步清空、collect/shadow 不缓存 | ✅ 完成 | T-X3 | | T-X6 | 模型池能力位过滤(外部采纳 cortiq capabilities):条目 capabilities{vision,tools,context_window}(缺省全兼容)+ filter_by_capabilities 硬过滤 + 代理请求需求推断(图片→vision/tools/上下文)与重定向 | ✅ 完成 | T-X6 | | T-X4 | 可解释决策留痕(外部采纳 ai-model-router 决策三元):sense_decisions 表(reasons/candidate_scores/rejected)+ grader 全模式落库 + /proxy/admin/sense-decisions 查询;附带修复 settings_provider 未注入导致 sense live 分流失效的潜伏 bug | ✅ 完成 | T-X4 | +| T-X5 | 云端评判晋升表×conformal 双闸门(外部采纳 cortiq promotion):PromotionTable 状态机(n_min/通过率/soak 晋升、退化降级、T3 一票否决)+ grader live 模式 T2→T1 应用 + 人工核验 verdict 回填(promo 标签) | ✅ 完成 | T-X5 | +| T-X7~X9 | P2 立项不实施:多协议透传(/v1/messages)/ 工具输出清洗沙箱(受限 DSL)/ 网关内置 A/B 实验框架 | 📋 立项 | 待排期 |