mirror of
https://github.com/ZhuLinsen/daily_stock_analysis
synced 2026-09-20 10:53:33 +08:00
fix: keep legacy alert triggers unattributed
This commit is contained in:
@@ -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 边界
|
||||
|
||||
|
||||
@@ -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)"
|
||||
|
||||
@@ -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(
|
||||
|
||||
Reference in New Issue
Block a user