"""应用装配:把规则加载、检测引擎、事件写入、收敛、通知与总线串起来。 DetectorApp 同时承载 REST API 所需的各 store 引用。 """ from __future__ import annotations from typing import List, Optional from .config import Config from .converger import AlertConverger from .engine import DetectionEngine from .models import MetricRule, MetricSample, RuleScope from .notifier import Notifier from .rule_loader import RuleLoader from .storage import ( AlertStore, Cache, EventStore, InMemoryMessageBus, MessageBus, RuleStore, create_alert_store, create_bus, create_cache, create_event_store, create_rule_store, ) class DetectorApp: def __init__(self, config: Optional[Config] = None) -> None: self.config = config or Config.from_env() self.rule_store: RuleStore = create_rule_store(self.config) self.event_store: EventStore = create_event_store(self.config) self.alert_store: AlertStore = create_alert_store(self.config) self.cache: Cache = create_cache(self.config) self.bus: MessageBus = create_bus(self.config) self.rule_loader = RuleLoader(self.rule_store, self.cache, self.config.rule_reload_interval) self.notifier = Notifier.memory() self.converger = AlertConverger( alert_store=self.alert_store, cache=self.cache, bus=self.bus, dedup_window=self.config.dedup_window, aggregate_window=self.config.aggregate_window, alerts_topic=self.config.alerts_topic, notifier=self.notifier, ) self.engine = DetectionEngine( rule_loader=self.rule_loader, evaluation_interval=self.config.evaluation_interval, recover_duration=self.config.recover_duration, on_event=self._on_event, ) def _on_event(self, event) -> None: # 事件先落库,再做收敛(收敛可能因去重被抑制)。 self.event_store.insert(event) self.converger.handle_event(event) # ------------------------------------------------------------------ # def start(self) -> None: self.rule_loader.load() self.rule_loader.start() def stop(self) -> None: self.rule_loader.stop() def seed_demo_rules(self) -> List[MetricRule]: """写入示例规则,便于本地联调/演示。""" demos = [ MetricRule( rule_id="r-cpu-high", name="CPU 使用率过高", metric="cpu_usage", aggregation="avg", operator="gt", threshold=90.0, for_duration="60s", severity="critical", scope=RuleScope(scope_type="all"), notify_channels=["email", "webhook"], enabled=True, ), MetricRule( rule_id="r-mem-high", name="内存使用率过高", metric="mem_used_percent", aggregation="avg", operator="gt", threshold=85.0, for_duration="60s", severity="warning", scope=RuleScope(scope_type="host_group", host_group="web"), notify_channels=["email"], enabled=True, ), MetricRule( rule_id="r-disk-full", name="磁盘使用率过高", metric="disk_used_percent", aggregation="avg", operator="gt", threshold=90.0, for_duration="30s", severity="critical", scope=RuleScope(scope_type="all"), notify_channels=["webhook"], enabled=True, ), ] created = [] for rule in demos: if self.rule_store.get_rule(rule.rule_id) is None: created.append(self.rule_store.create_rule(rule)) self.rule_loader.load() return created # ------------------------------------------------------------------ # def ingest(self, host_id: str, host_group: Optional[str], samples: List[MetricSample]): """接收一批样本并驱动检测,返回触发的事件。""" return self.engine.handle_batch(samples, host_id, host_group)