125 lines
4.3 KiB
Python

"""应用装配:把规则加载、检测引擎、事件写入、收敛、通知与总线串起来。
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)