diff --git a/docs/alerts.md b/docs/alerts.md index f41db0d39..482ae0bf8 100644 --- a/docs/alerts.md +++ b/docs/alerts.md @@ -440,7 +440,7 @@ worker 会把 `triggered`、`skipped`、`degraded`、`failed` 写入 `alert_trig - 触发归属使用 `rule_id + rule_lifecycle_id`:每次创建规则都会生成不可复用的 lifecycle ID,worker 在加载规则时固定该 ID,并随触发记录一起写入。即使旧 worker 已经开始计算、期间规则被删除且 SQLite 复用了整数 ID,迟到的旧触发也不会归到新规则。`rule_id` 为空的记录计入 `unattributed_trigger_count`;引用已不存在规则或 lifecycle 不匹配的记录计入 `orphaned_trigger_count`。 - 摘要的总数、状态分组、规则分组和孤儿计数在同一个 SQLite 读事务快照内完成,后台 worker 同时写入触发记录时不会返回“总数小于状态分组之和”的自相矛盾结果。 -该概览是 Issue #2281 的基础阶段,只解决全局统计与可靠归属。它不引入“论文失效”规则类型,也不把普通告警命名为失效信号;`impact`、`affected_entities` 和新的 thesis invalidation 规则语义留待后续独立 PR。启动时会为既有 SQLite 的 `alert_rules` / `alert_triggers` 增加 lifecycle 列,并只把能按原创建时间规则可靠匹配的既有触发回填到当前 lifecycle;无法可靠匹配的历史保留为 orphan。回滚代码不会自动删除新增列或历史数据,旧版本会忽略这些附加列。 +该概览是 Issue #2281 的基础阶段,只解决全局统计与可靠归属。它不引入“论文失效”规则类型,也不把普通告警命名为失效信号;`impact`、`affected_entities` 和新的 thesis invalidation 规则语义留待后续独立 PR。启动时会为既有 SQLite 的 `alert_rules` / `alert_triggers` 增加 lifecycle 列;既有规则会获得 lifecycle,但迁移前的触发记录无法证明属于哪一代规则,因此不按时间戳猜测回填,统一保留为 orphan。多实例同时首次升级时,重复加列会安全收敛。回滚代码不会自动删除新增列或历史数据,旧版本会忽略这些附加列。 ## Phase 边界 diff --git a/src/storage.py b/src/storage.py index 7aac2f2ae..18217df9a 100644 --- a/src/storage.py +++ b/src/storage.py @@ -1487,27 +1487,36 @@ class DatabaseManager(metaclass=_DatabaseManagerMeta): trigger_columns = { column["name"] for column in inspector.get_columns(AlertTriggerRecord.__tablename__) } + for table, column, column_type, existing in ( + ("alert_rules", "lifecycle_id", "VARCHAR(36)", rule_columns), + ("alert_triggers", "rule_lifecycle_id", "VARCHAR(36)", trigger_columns), + ): + if column in existing: + continue + for attempt in range(self._sqlite_write_retry_max + 1): + try: + with self._engine.begin() as connection: + connection.exec_driver_sql( + f"ALTER TABLE {table} ADD COLUMN {column} {column_type}" + ) + existing.add(column) + break + except OperationalError as exc: + if self._is_sqlite_duplicate_column_error(exc, column): + existing.add(column) + break + if self._is_sqlite_locked_error(exc) and attempt < self._sqlite_write_retry_max: + delay = self._sqlite_write_retry_base_delay * (2 ** attempt) + if delay > 0: + time.sleep(delay) + continue + raise + with self._engine.begin() as connection: - if "lifecycle_id" not in rule_columns: - connection.exec_driver_sql( - "ALTER TABLE alert_rules ADD COLUMN lifecycle_id VARCHAR(36)" - ) - if "rule_lifecycle_id" not in trigger_columns: - connection.exec_driver_sql( - "ALTER TABLE alert_triggers ADD COLUMN rule_lifecycle_id VARCHAR(36)" - ) connection.exec_driver_sql( "UPDATE alert_rules SET lifecycle_id = lower(hex(randomblob(16))) " "WHERE lifecycle_id IS NULL OR lifecycle_id = ''" ) - connection.exec_driver_sql( - "UPDATE alert_triggers SET rule_lifecycle_id = (" - "SELECT lifecycle_id FROM alert_rules " - "WHERE alert_rules.id = alert_triggers.rule_id " - "AND alert_triggers.triggered_at IS NOT NULL " - "AND alert_triggers.triggered_at >= alert_rules.created_at" - ") WHERE rule_lifecycle_id IS NULL" - ) connection.exec_driver_sql( "CREATE INDEX IF NOT EXISTS ix_alert_rules_lifecycle_id " "ON alert_rules (lifecycle_id)" diff --git a/tests/test_alert_api.py b/tests/test_alert_api.py index 776d71001..e8d4ee9e1 100644 --- a/tests/test_alert_api.py +++ b/tests/test_alert_api.py @@ -932,6 +932,25 @@ class AlertApiTestCase(unittest.TestCase): self.assertEqual(by_rule_id[replacement_rule["id"]]["trigger_count"], 0) self.assertEqual(payload["orphaned_trigger_count"], 1) + def test_monitor_summary_keeps_lifecycle_less_legacy_trigger_orphaned(self) -> None: + rule = self._create_rule({"name": "current", "target": "600519"}) + with self.db.get_session() as session: + session.add( + AlertTriggerRecord( + rule_id=rule["id"], + rule_lifecycle_id=None, + target="600519", + status="triggered", + ) + ) + session.commit() + + payload = self.client.get("/api/v1/alerts/summary").json() + + by_rule_id = {item["rule_id"]: item for item in payload["rules"]} + self.assertEqual(by_rule_id[rule["id"]]["trigger_count"], 0) + self.assertEqual(payload["orphaned_trigger_count"], 1) + def test_monitor_summary_uses_one_read_snapshot_during_concurrent_trigger_write(self) -> None: rule = self._create_rule({"name": "snapshot", "target": "600519"}) AlertRepository(self.db).create_trigger(