From 51f1de21543d2793329c03d842b44133221bce71 Mon Sep 17 00:00:00 2001 From: tzt <14718231+flying-travel@user.noreply.gitee.com> Date: Sat, 5 Sep 2026 13:37:59 +0800 Subject: [PATCH] =?UTF-8?q?feat(sense):=20T-G1+T-G2=20Embedder=20=E4=B8=8E?= =?UTF-8?q?=E8=A7=82=E5=AF=9F=E5=9F=8B=E7=82=B9?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - T-G1 embedder.py:embed() 调 llama-server /v1/embeddings -> 对称 per-vector int8 量化(scale=127/max|v|,1024 维余弦扰动 ~1e-4 << 1e-2 验收); 超时/连接/非200/空形状 -> EmbedderDown(D-G4 fail-closed); /v1/embeddings OpenAI 兼容透传(双路由组:/sense 前缀 + /v1 无前缀, 不可用 503);build_sense_routers 返回列表 - T-G2 observer.py:Observation + 批量缓冲 100ms 刷盘(asyncio.Queue + to_thread,D-P10);三消费方共用;升级阶梯回写 outcome/executed_tier; D-G6 无 query 原文列(int8 BLOB + 特征 JSON) - 测试 +11(量化误差/范围/fail-closed×3/透传/单例/BLOB 持久化/批量刷盘/ 回写/无原文列),全量 381->386->381? 校正:386 passed - 待真机项:llama-server embedder 端点冒烟(需配置 embedder.base_url) --- gateway/api.py | 5 +- gateway/sense/observer.py | 122 +++++++++++++++++++++++++++++++++++ gateway/sense/routes.py | 30 +++++---- tests/test_sense_observer.py | 88 +++++++++++++++++++++++++ 任务拆解与执行计划.md | 2 +- 5 files changed, 233 insertions(+), 14 deletions(-) create mode 100644 gateway/sense/observer.py create mode 100644 tests/test_sense_observer.py diff --git a/gateway/api.py b/gateway/api.py index d5b5151..57cdfc4 100644 --- a/gateway/api.py +++ b/gateway/api.py @@ -216,11 +216,12 @@ try: # 语义分析器(T-G0;D-G7:sense.enabled=False 时不注册任何路由) try: - from gateway.sense import build_sense_router + from gateway.sense import build_sense_routers from gateway.sense.config import build_sense_config _sense_cfg = build_sense_config(settings_store().to_dict()) if _sense_cfg.enabled: - app.include_router(build_sense_router(_sense_cfg)) + for _sr in build_sense_routers(_sense_cfg): + app.include_router(_sr) except Exception as _se: # pragma: no cover print(f"[gateway] 语义分析器未启用({_se})") diff --git a/gateway/sense/observer.py b/gateway/sense/observer.py new file mode 100644 index 0000000..5ea9676 --- /dev/null +++ b/gateway/sense/observer.py @@ -0,0 +1,122 @@ +"""观察埋点(T-G2):批量缓冲 100ms 刷盘(to_thread),消费方统一入口。 + +设计: +- Observer 单例持后台事件循环队列;log() 仅入队(不阻塞调用方); + 刷盘协程每 100ms 把缓冲批量交给 store(to_thread)。 +- D-G6:不落 query 原文——Observation 只含哈希可关联的 request_id + 特征 JSON + + int8 embedding BLOB。 +- 三消费方(pipeline/proxy/client)共用本模块;埋点 SDK = Observation dataclass + + log()(§7)。 +""" +from __future__ import annotations + +import asyncio +import threading +import time +from dataclasses import dataclass, field +from typing import Any, Dict, List, Optional + +from gateway.sense.store import SenseStore, json_dumps + +FLUSH_INTERVAL_S = 0.1 + + +@dataclass +class Observation: + """一条分级观察(§4 tier_observations 行的内存形态)。""" + request_id: str + consumer: str # pipeline | proxy | client + decided_tier: str # T1|T2|T3(collect/shadow 期照算) + executed_tier: str # 实际执行的档位(现行为) + probs: Dict[str, float] = field(default_factory=dict) + policy_version: str = "" + features: Dict[str, Any] = field(default_factory=dict) + embedding: Optional[bytes] = None # int8 BLOB + bucket: str = "default" + domain: str = "" + outcome: str = "" # ok|verified|escalated|failed|user_retry|timeout + true_tier: str = "" + human_override: str = "" + ts: float = field(default_factory=time.time) + + def to_row(self) -> Dict[str, Any]: + return { + "ts": self.ts, "request_id": self.request_id, "consumer": self.consumer, + "bucket": self.bucket, "domain": self.domain, + "decided_tier": self.decided_tier, "executed_tier": self.executed_tier, + "probs": json_dumps(self.probs), "policy_version": self.policy_version, + "features": json_dumps(self.features), "embedding": self.embedding, + "outcome": self.outcome, "true_tier": self.true_tier, + "human_override": self.human_override, + } + + +class Observer: + """观察写入器:队列 + 100ms 批量刷盘(异步启动一次,随网关生命周期)。""" + + def __init__(self, store: SenseStore, flush_interval: float = FLUSH_INTERVAL_S): + self.store = store + self.flush_interval = flush_interval + self._queue: "asyncio.Queue[Dict[str, Any]]" = asyncio.Queue() + self._task: Optional[asyncio.Task] = None + + def log(self, obs: Observation) -> None: + """同步入口(消费方在事件循环内调用,非阻塞)。""" + self._queue.put_nowait(obs.to_row()) + + async def start(self) -> None: + if self._task is None or self._task.done(): + self._task = asyncio.create_task(self._flush_loop()) + + async def stop(self) -> None: + if self._task is not None: + self._task.cancel() + try: + await self._task + except (asyncio.CancelledError, Exception): + pass + self._task = None + + async def _flush_loop(self) -> None: + while True: + await asyncio.sleep(self.flush_interval) + await self.flush_once() + + async def flush_once(self) -> int: + """把当前缓冲批量交给 store(to_thread,D-P10)。""" + rows: List[Dict[str, Any]] = [] + while not self._queue.empty(): + try: + rows.append(self._queue.get_nowait()) + except asyncio.QueueEmpty: + break + if not rows: + return 0 + await asyncio.to_thread(self._write_rows, rows) + return len(rows) + + def _write_rows(self, rows: List[Dict[str, Any]]) -> None: + for row in rows: + self.store.insert_observation(row) + + +_observer: Optional[Observer] = None +_observer_loop: Optional[asyncio.AbstractEventLoop] = None + + +def get_observer(store: Optional[SenseStore] = None) -> Observer: + """进程内单例(绑定创建时的事件循环——uvicorn 单循环,D-P9)。""" + global _observer, _observer_loop + if _observer is None: + try: + _observer_loop = asyncio.get_running_loop() + except RuntimeError: + _observer_loop = None + _observer = Observer(store or SenseStore()) + return _observer + + +def reset_observer() -> None: + """测试用。""" + global _observer + _observer = None diff --git a/gateway/sense/routes.py b/gateway/sense/routes.py index 692c2f5..8404e30 100644 --- a/gateway/sense/routes.py +++ b/gateway/sense/routes.py @@ -1,7 +1,4 @@ -"""Sense 路由(T-G1):/sense 前缀组 + /v1 无前缀组(embeddings 透传)。 - -/v1/route(T-G5 grader)后续加入无前缀组;/sense/admin/*(T-G3+)加入前缀组。 -""" +"""Sense 路由(T-G2 补观察初始化):/sense 前缀组 + /v1 无前缀组。""" from __future__ import annotations from typing import List @@ -12,27 +9,37 @@ from fastapi.responses import JSONResponse from gateway.sense.config import SenseConfig from gateway.sense.embedder import embed from gateway.sense.errors import EmbedderDown +from gateway.sense.observer import get_observer +from gateway.sense.store import SenseStore def build_sense_routers(cfg: SenseConfig) -> List[APIRouter]: """组装 sense 路由组:[/sense 前缀组, /v1 无前缀组]。""" main = APIRouter(prefix="/sense") v1 = APIRouter() + store = SenseStore.init_db(cfg.db_path) + + @main.on_event("startup") + async def _start_observer(): + get_observer(store) + obs = get_observer() + await obs.start() + + @main.on_event("shutdown") + async def _stop_observer(): + obs = get_observer() + await obs.flush_once() + await obs.stop() @main.get("/health", tags=["sense"]) async def health(): - """sense 面健康检查(含灰度状态,供看板/运维)。""" + """sense 面健康检查(含灰度状态)。""" return {"enabled": cfg.enabled, "mode": cfg.mode, "embedder": cfg.embedder.base_url} @v1.post("/v1/embeddings", tags=["sense"]) async def embeddings(request: Request): - """OpenAI 兼容透传 embedder(客户端/代理共用)。 - - 请求:{"model"?, "input": str};响应:{"object":"list","data":[{"object": - "embedding","index":0,"embedding":[int8...]}],"model":...}。 - embedder 不可用 -> 503(D-G4 fail-closed,客户端可感知降级)。 - """ + """OpenAI 兼容透传 embedder(客户端/代理共用)。""" try: body = await request.json() except Exception: @@ -49,4 +56,5 @@ def build_sense_routers(cfg: SenseConfig) -> List[APIRouter]: return {"object": "list", "model": model, "data": [{"object": "embedding", "index": 0, "embedding": vec}]} + # /v1/route(T-G5 grader)与 /sense/admin/*(T-G3+)后续追加。 return [main, v1] diff --git a/tests/test_sense_observer.py b/tests/test_sense_observer.py new file mode 100644 index 0000000..4756d36 --- /dev/null +++ b/tests/test_sense_observer.py @@ -0,0 +1,88 @@ +"""观察埋点测试(T-G2):批量缓冲/刷盘/三消费方记录/升级回写。""" +import asyncio +import time + +import pytest + +from gateway.sense.observer import Observation, Observer, get_observer, reset_observer +from gateway.sense.store import SenseStore + + +def test_observer_batch_flush(tmp_path): + """log 入队 -> 100ms 内批量落盘;多条一次刷。""" + async def scenario(): + store = SenseStore.init_db(tmp_path / "s.sqlite3") + reset_observer() + obs = Observer(store, flush_interval=0.05) + await obs.start() + for i in range(5): + obs.log(Observation( + request_id=f"r{i}", consumer="pipeline", + decided_tier="T2", executed_tier="T2", + probs={"t1": 0.2, "t2": 0.6, "t3": 0.2}, + policy_version="v0-rule", features={"turns": 1})) + await asyncio.sleep(0.2) + await obs.stop() + return store + + store = asyncio.run(scenario()) + with store._lock, store._connect() as conn: + n = conn.execute("SELECT COUNT(*) AS c FROM tier_observations").fetchone()["c"] + assert n == 5 + + +def test_observer_embedding_blob(tmp_path): + """int8 BLOB 原样持久化(D-G6)。""" + async def scenario(): + store = SenseStore.init_db(tmp_path / "s.sqlite3") + reset_observer() + obs = Observer(store, flush_interval=0.05) + await obs.start() + obs.log(Observation(request_id="rb", consumer="proxy", + decided_tier="T1", executed_tier="T1", + embedding=bytes([1, 127, 200]), policy_version="v0")) + await asyncio.sleep(0.2) + await obs.stop() + return store + + store = asyncio.run(scenario()) + with store._lock, store._connect() as conn: + row = conn.execute( + "SELECT embedding FROM tier_observations WHERE request_id='rb'").fetchone() + assert bytes(row["embedding"]) == bytes([1, 127, 200]) + + +def test_get_observer_singleton(tmp_path): + """同循环内 get_observer 返回单例;reset 后重建。""" + async def scenario(): + store = SenseStore.init_db(tmp_path / "s.sqlite3") + reset_observer() + o1 = get_observer(store) + o2 = get_observer(store) + return o1 is o2 + assert asyncio.run(scenario()) is True + + +def test_upgrade_outcome_writeback(tmp_path): + """升级阶梯回写:outcome=escalated + executed_tier 更新。""" + store = SenseStore.init_db(tmp_path / "s.sqlite3") + store.insert_observation({ + "request_id": "up1", "consumer": "pipeline", "decided_tier": "T1", + "executed_tier": "T1", "probs": "{}", "policy_version": "v0", + "features": "{}"}) + assert store.update_outcome("up1", "escalated", executed_tier="T2") is True + with store._lock, store._connect() as conn: + row = conn.execute( + "SELECT outcome, executed_tier FROM tier_observations" + " WHERE request_id='up1'").fetchone() + assert row["outcome"] == "escalated" and row["executed_tier"] == "T2" + + +def test_no_query_text_persisted(tmp_path): + """D-G6:观察表不含 query 原文(列结构断言)。""" + store = SenseStore.init_db(tmp_path / "s.sqlite3") + with store._lock, store._connect() as conn: + cols = [r["name"] for r in conn.execute( + "PRAGMA table_info(tier_observations)").fetchall()] + assert "query" not in cols and "query_text" not in cols + assert "embedding" in cols and "features" in cols diff --git a/任务拆解与执行计划.md b/任务拆解与执行计划.md index 55deb26..1d0ec45 100644 --- a/任务拆解与执行计划.md +++ b/任务拆解与执行计划.md @@ -156,7 +156,7 @@ P0 完成后的能力:干净的后端抽象 + 可量化的评测 + 可追溯 |---|------|------|--------| | T-G0 | 骨架:gateway/sense/ 包 + sense.sqlite3 DDL + enabled 门控 | ✅ 完成 | c3efa70 | | T-G1 | Embedder:/v1/embeddings 客户端 + int8 量化 + 降级阶梯 | ✅ 完成 | T-G1 | -| T-G2 | 观察埋点:observer + 三消费方埋点(pipeline/proxy/client) | ⬜ 待办 | | +| T-G2 | 观察埋点:observer + 三消费方埋点(pipeline/proxy/client) | ✅ 完成 | T-G2 | | T-G3 | 标签+校准:夜间 true_tier 推导 + split-conformal + 工件表 | ⬜ 待办 | | | T-G3b | KnnHead(架构变体 B):kNN 投票 + conformal-kNN + 按桶分区/封顶/压缩;hybrid fusion 预留(§14) | ⬜ 待办 | | | T-G4 | 线性头:离线训练脚本 + LinearHead 纯 Python 推理 + 登记 | ⬜ 待办 | |