fix: keep the postgres backfill from disabling telemetry or locking the table
This commit is contained in:
@@ -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()
|
||||
|
||||
@@ -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)。"""
|
||||
|
||||
Reference in New Issue
Block a user