103 lines
3.3 KiB
Python
103 lines
3.3 KiB
Python
from typing import List
|
||
|
||
from sqlalchemy import insert, select, func
|
||
from sqlalchemy.orm import Session
|
||
|
||
from models.log_alert import LogAlert
|
||
from schemas.log_alert import LogAlertTrigger, LogAlertResponse
|
||
from service import openobserve_service
|
||
from service.message_service import send_wechat_message, MessageEnum
|
||
from utils.time import get_timestamp_range, format_timestamp
|
||
from config.logging import logger
|
||
|
||
|
||
def trigger_alert(log_alter: LogAlertTrigger, db: Session) -> bool:
|
||
logger.info(f"{format_timestamp(log_alter.timestamp)} 发生告警:{log_alter.alter_name}")
|
||
|
||
start_time, end_time = get_timestamp_range(log_alter.timestamp, 60)
|
||
result = openobserve_service.get_log(start_time, end_time)
|
||
hits = result.get("hits")
|
||
|
||
# 判断日志是否为空
|
||
if len(hits) == 0:
|
||
return False
|
||
|
||
# 判断日志告警是否已经存在
|
||
stmt = select(LogAlert).where(LogAlert.alter_time == format_timestamp(log_alter.timestamp))
|
||
log_db_alter = db.execute(stmt).first()
|
||
|
||
if log_db_alter is not None:
|
||
return False
|
||
|
||
# 批量插入
|
||
batch_insert(log_alter, hits, db)
|
||
|
||
# 通知消息
|
||
# notify_message(log_alter)
|
||
|
||
return True
|
||
|
||
|
||
def batch_insert(log_alter: LogAlertTrigger, hits: [], db: Session) -> bool:
|
||
stmt = insert(LogAlert).values(
|
||
[
|
||
{
|
||
"alter_name": log_alter.alter_name,
|
||
"alter_time": format_timestamp(log_alter.timestamp),
|
||
"log_time": format_timestamp(hit["_timestamp"]),
|
||
"log_topic": hit["topic"],
|
||
"log_message": hit["message"],
|
||
"log_status": "pending"
|
||
}
|
||
for hit in hits
|
||
]
|
||
)
|
||
|
||
db.execute(stmt)
|
||
db.commit()
|
||
|
||
return True
|
||
|
||
|
||
def notify_message(log_alter: LogAlertTrigger):
|
||
# 企业微信通知
|
||
alter_message = f"🔔 **告警服务**:{log_alter.alter_name} \n 🕒 **告警时间**:{format_timestamp(log_alter.timestamp)}"
|
||
send_wechat_message(MessageEnum.MARKDOWN2, alter_message)
|
||
|
||
|
||
def query_alert(db: Session):
|
||
stmt = (
|
||
select(
|
||
LogAlert.alter_name.label("alterName"),
|
||
func.strftime("%Y-%m-%d %H:%M:%S", LogAlert.alter_time).label("alterTime"),
|
||
func.strftime("%Y-%m-%d %H:%M:%S", LogAlert.log_time).label("logTime"),
|
||
LogAlert.log_topic.label("logTopic"),
|
||
LogAlert.log_message.label("logMessage"),
|
||
LogAlert.log_status.label("logStatus")
|
||
)
|
||
.order_by(LogAlert.alter_name.desc(), LogAlert.alter_time.desc())
|
||
)
|
||
|
||
results = db.execute(stmt).fetchall()
|
||
|
||
return [LogAlertResponse.model_validate(result) for result in results]
|
||
|
||
|
||
def query_alert_by_name(name: str, db: Session) -> List[LogAlertResponse]:
|
||
stmt = (
|
||
select(
|
||
LogAlert.alter_name.label("alterName"),
|
||
func.strftime("%Y-%m-%d %H:%M:%S", LogAlert.alter_time).label("alterTime"),
|
||
func.strftime("%Y-%m-%d %H:%M:%S", LogAlert.log_time).label("logTime"),
|
||
LogAlert.log_topic.label("logTopic"),
|
||
LogAlert.log_message.label("logMessage"),
|
||
LogAlert.log_status.label("logStatus")
|
||
)
|
||
.where(LogAlert.alter_name == name)
|
||
.order_by(LogAlert.alter_time.desc(), LogAlert.log_time.asc())
|
||
)
|
||
|
||
results = db.execute(stmt).fetchall()
|
||
|
||
return [LogAlertResponse.model_validate(result) for result in results]
|