diff --git a/src/polygateway/pricing.py b/src/polygateway/pricing.py index 0b4d249..0e394dc 100644 --- a/src/polygateway/pricing.py +++ b/src/polygateway/pricing.py @@ -47,6 +47,8 @@ class PricingTable: def __init__(self, prices: Mapping[str, ModelPrice]) -> None: self._prices = dict(prices) self._warned: set[str] = set() + # 独立集合: 与"未知 model"的告警去重键分开,避免 model 名恰好撞上时互相抑制 + self._warned_clamp: set[str] = set() @classmethod def from_file(cls, path: Path | str) -> PricingTable: @@ -111,9 +113,8 @@ class PricingTable: """命中数按输入总数夹取: 网关口径异常不得算出负成本(每 model 只警告一次)。""" if cached <= prompt_tokens: return cached - key = f"{model}:cached_over_prompt" - if key not in self._warned: - self._warned.add(key) + if model not in self._warned_clamp: + self._warned_clamp.add(model) logger.warning( "model {!r} 上报的缓存命中 {} 超过输入总数 {},按总数夹取计价", model, diff --git a/src/polygateway/telemetry/postgres.py b/src/polygateway/telemetry/postgres.py index efb6469..f8777f2 100644 --- a/src/polygateway/telemetry/postgres.py +++ b/src/polygateway/telemetry/postgres.py @@ -47,8 +47,14 @@ CREATE TABLE IF NOT EXISTS llm_calls ( # 新列排在 created_at 之后: 与旧表 ALTER 追加的位置一致(见 sqlite.py 同款注释) _BACKFILL = ( - "ALTER TABLE llm_calls ADD COLUMN IF NOT EXISTS cached_prompt_tokens INTEGER", - "ALTER TABLE llm_calls ADD COLUMN IF NOT EXISTS model_reported TEXT", + ("cached_prompt_tokens", "ALTER TABLE llm_calls ADD COLUMN cached_prompt_tokens INTEGER"), + ("model_reported", "ALTER TABLE llm_calls ADD COLUMN model_reported TEXT"), +) + +# 探测现有列;尊重 search_path(to_regclass 按当前 search_path 解析) +_EXISTING_COLUMNS = ( + "SELECT attname FROM pg_attribute " + "WHERE attrelid = to_regclass('llm_calls') AND attnum > 0 AND NOT attisdropped" ) _COLUMNS = ( @@ -114,9 +120,7 @@ class PostgresRecorder: self._pool = await asyncpg.create_pool(self._dsn, timeout=10) async with self._pool.acquire() as conn: await conn.execute(_DDL) - for statement in _BACKFILL: - # 已存在的旧表补新列(issue #3);ADD COLUMN IF NOT EXISTS 原生幂等 - await conn.execute(statement) + await self._backfill_columns(conn) self._schema_ready = True return self._pool except asyncio.CancelledError: @@ -126,6 +130,29 @@ class PostgresRecorder: logger.warning("Postgres 遥测初始化失败,后续记录降级为 no-op: {}", exc) return None + async def _backfill_columns(self, conn: object) -> None: + """给已存在的旧表补新列(issue #3);**先探测再 ALTER,失败绝不置 `_failed`**。 + + 两条纪律各有实测理由: + ① 不置 `_failed`: 应用账号只有 INSERT 权限时,`ALTER TABLE` 的 ownership + 检查早于 `IF NOT EXISTS` 的存在性判断——列明明齐全也会失败。置位会让 + 整个 recorder 永久 no-op,与「补列失败只降级为逐行丢弃」的承诺相悖 + (SQLite 侧同款守卫,两侧必须对称)。 + ② 先探测: `ADD COLUMN IF NOT EXISTS` 即便列已存在,也会**先取 ACCESS + EXCLUSIVE 锁**再判存在性(实测会被一个开着的读事务阻塞)。遥测是内联 + await,让每个进程的首次写入都去抢共享审计表的排他锁,等于用记录基础设施 + 拖垮业务调用。探测走 ACCESS SHARE,稳态下一条 ALTER 都不会发。 + """ + try: + existing = {row["attname"] for row in await conn.fetch(_EXISTING_COLUMNS)} # type: ignore[attr-defined] + for column, statement in _BACKFILL: + if column not in existing: + await conn.execute(statement) # type: ignore[attr-defined] + except asyncio.CancelledError: + raise + except Exception as exc: + logger.warning("Postgres 遥测补列失败(写入将逐行降级): {}", exc) + async def record_llm_call(self, **fields: object) -> None: """写一行遥测;单条失败逐条 warning 丢弃(两级降级之二),绝不冒泡。""" pool = await self._ensure_ready() diff --git a/src/polygateway/telemetry/sqlite.py b/src/polygateway/telemetry/sqlite.py index fb992fc..91eadb0 100644 --- a/src/polygateway/telemetry/sqlite.py +++ b/src/polygateway/telemetry/sqlite.py @@ -104,15 +104,20 @@ class SQLiteRecorder: return try: existing = {row[1] for row in self._conn.execute("PRAGMA table_info(llm_calls)")} - for column, decl in _BACKFILL_COLUMNS: - if column in existing: - continue - self._conn.execute(f"ALTER TABLE llm_calls ADD COLUMN {column} {decl}") - self._conn.commit() except sqlite3.Error as exc: - # duplicate column: 多进程共库时后到者必然撞上,属预期竞态,视为成功 - if "duplicate column" not in str(exc).lower(): - logger.warning("SQLite 遥测补列失败(写入将逐行降级): {}", exc) + logger.warning("SQLite 遥测列探测失败(写入将逐行降级): {}", exc) + return + for column, decl in _BACKFILL_COLUMNS: + if column in existing: + continue + # 逐列独立 try: 一列撞上 duplicate 不得让后面的列漏补 + try: + self._conn.execute(f"ALTER TABLE llm_calls ADD COLUMN {column} {decl}") + self._conn.commit() + except sqlite3.Error as exc: + # duplicate column: 多进程共库时后到者必然撞上,属预期竞态,视为成功 + if "duplicate column" not in str(exc).lower(): + logger.warning("SQLite 遥测补列失败(写入将逐行降级): {}", exc) async def record_llm_call(self, **fields: object) -> None: """写一行遥测;字段集合即 20 字段冻结签名(ports.TelemetryRecorder)。""" diff --git a/src/polygateway/transports/openai_compat.py b/src/polygateway/transports/openai_compat.py index 7703f54..f0164e9 100644 --- a/src/polygateway/transports/openai_compat.py +++ b/src/polygateway/transports/openai_compat.py @@ -178,10 +178,13 @@ def _coerce_cached_tokens(usage: Any) -> int | None: def _coerce_model_reported(value: Any) -> str | None: - """取响应体的 model 字段(issue #3);非 str 或空白串一律 None。""" + """取响应体的 model 字段(issue #3);非 str 或空白串一律 None,收口时去空白。 + + 去空白不是洁癖: 下游拿这个串做实验快照的 key,`" m "` 与 `"m"` 会造成假分叉。 + """ if not isinstance(value, str) or not value.strip(): return None - return value + return value.strip() def _resolve_stream_usage(sink: dict[str, Any], salvaged: bool) -> tuple[int, int, str]: diff --git a/tests/unit/test_telemetry.py b/tests/unit/test_telemetry.py index 390c606..e3aa4ed 100644 --- a/tests/unit/test_telemetry.py +++ b/tests/unit/test_telemetry.py @@ -219,6 +219,78 @@ class TestSQLiteColumnBackfill: recorder.close() +class _FakePgConn: + """记录执行过的语句;可让 ALTER 抛错以模拟权限不足。""" + + def __init__(self, existing: list[str], *, fail_alter: bool = False): + self.existing = existing + self.fail_alter = fail_alter + self.statements: list[str] = [] + + async def execute(self, sql, *args): + self.statements.append(sql) + if sql.startswith("ALTER TABLE") and self.fail_alter: + raise RuntimeError("must be owner of table llm_calls") + + async def fetch(self, sql, *args): + self.statements.append(sql) + return [{"attname": name} for name in self.existing] + + +class _FakePgPool: + def __init__(self, conn): + self._conn = conn + + def acquire(self): + conn = self._conn + + class _Ctx: + async def __aenter__(self): + return conn + + async def __aexit__(self, *exc): + return False + + return _Ctx() + + +class TestPostgresBackfillDiscipline: + """PG 补列必须与 SQLite 侧对称: 失败只逐行降级,且稳态不抢排他锁(issue #3)。""" + + _LEGACY = ["call_id", "cost", "created_at"] + _CURRENT = ["call_id", "cost", "created_at", "cached_prompt_tokens", "model_reported"] + + def _recorder(self, conn): + from polygateway.telemetry.postgres import PostgresRecorder + + return PostgresRecorder("postgresql://u:p@h:5432/polygateway", pool=_FakePgPool(conn)) + + async def test_alter_failure_does_not_disable_the_recorder(self): + """ALTER 失败(如账号只有 INSERT 权限)不得置 _failed —— 那会让遥测全灭。""" + conn = _FakePgConn(self._LEGACY, fail_alter=True) + recorder = self._recorder(conn) + await _record_minimal(recorder) # 不得抛 + assert recorder._failed is False + assert any(s.startswith("INSERT INTO llm_calls") for s in conn.statements) + + async def test_no_alter_when_columns_already_exist(self): + """ADD COLUMN IF NOT EXISTS 即使列已存在也会先抢 ACCESS EXCLUSIVE 锁, + + 而遥测是内联 await——稳态下必须一条 ALTER 都不发,否则每个进程的首次 + 写入都会去锁共享审计表。 + """ + conn = _FakePgConn(self._CURRENT) + await _record_minimal(self._recorder(conn)) + assert not [s for s in conn.statements if s.startswith("ALTER TABLE")] + + async def test_missing_columns_are_added_once(self): + conn = _FakePgConn(self._LEGACY) + await _record_minimal(self._recorder(conn)) + altered = [s for s in conn.statements if s.startswith("ALTER TABLE")] + assert len(altered) == 2 + assert all("IF NOT EXISTS" not in s for s in altered) # 探测已确认缺列,无需再判 + + class _MemoryRecorder: def __init__(self): self.rows = []