Files
projectAIpopular/gateway/sense/observer.py
T
tzt 51f1de2154 feat(sense): T-G1+T-G2 Embedder 与观察埋点
- 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)
2026-09-05 13:37:59 +08:00

123 lines
4.4 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""观察埋点(T-G2):批量缓冲 100ms 刷盘(to_thread),消费方统一入口。
设计:
- Observer 单例持后台事件循环队列;log() 仅入队(不阻塞调用方);
刷盘协程每 100ms 把缓冲批量交给 storeto_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|T3collect/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:
"""把当前缓冲批量交给 storeto_threadD-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