"""观察埋点(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