diff --git a/.gitignore b/.gitignore index 7ddac70..3a2e874 100644 --- a/.gitignore +++ b/.gitignore @@ -65,4 +65,5 @@ media/ .DS_Store logs/ -packages/ \ No newline at end of file +packages/ +*.db \ No newline at end of file diff --git a/apis/__pycache__/log.cpython-310.pyc b/apis/__pycache__/log.cpython-310.pyc index 656dd7c..c05b3b1 100644 Binary files a/apis/__pycache__/log.cpython-310.pyc and b/apis/__pycache__/log.cpython-310.pyc differ diff --git a/apis/log.py b/apis/log.py index 7acd199..3784d7a 100644 --- a/apis/log.py +++ b/apis/log.py @@ -1,8 +1,10 @@ +from typing import List + from fastapi import APIRouter, Depends from sqlalchemy.orm import Session from config.database import get_db -from schemas.log_alert import LogAlertTrigger +from schemas.log_alert import LogAlertTrigger, LogAlertResponse from service import log_alert_service router = APIRouter( @@ -12,11 +14,16 @@ router = APIRouter( ) -@router.get("", summary="查询日志告警") +@router.get("", summary="查询所有日志告警") def query_alert(db: Session = Depends(get_db)): return log_alert_service.query_alert(db) +@router.get("/name", summary="查询指定日志告警", response_model=List[LogAlertResponse]) +def query_alert(name: str, db: Session = Depends(get_db)): + return log_alert_service.query_alert_by_name(name, db) + + @router.post("", summary="触发日志告警") def trigger_log_alter(log_alter: LogAlertTrigger, db: Session = Depends(get_db)) -> bool: return log_alert_service.trigger_alert(log_alter, db) diff --git a/config/__pycache__/database.cpython-310.pyc b/config/__pycache__/database.cpython-310.pyc index 0c1b3f5..5a84fe9 100644 Binary files a/config/__pycache__/database.cpython-310.pyc and b/config/__pycache__/database.cpython-310.pyc differ diff --git a/config/database.py b/config/database.py index 3440129..ce92b19 100644 --- a/config/database.py +++ b/config/database.py @@ -1,5 +1,4 @@ from sqlalchemy import create_engine -from sqlalchemy.ext.declarative import declarative_base from sqlalchemy.orm import sessionmaker from config.setting import settings diff --git a/config/logging.py b/config/logging.py new file mode 100644 index 0000000..eee43be --- /dev/null +++ b/config/logging.py @@ -0,0 +1,32 @@ +import sys + +from loguru import logger + +from config.setting import settings + +# 日志级别 +LOG_LEVEL = settings.log_level.upper() + +# 日志格式 +STDOUT_FORMAT = ( + "{time:YYYY-MM-DD HH:mm:ss.SSS} | " + "{level: <8} | " + "{name}:{function}:{line} - " + "{message}" +) + +# 移除默认处理器 +logger.remove() + +# 添加控制台处理器 +logger.add( + sink=sys.stdout, + level=LOG_LEVEL, + format=STDOUT_FORMAT, + colorize=True, + backtrace=True, # 显示完整异常堆栈 + diagnose=True, # 显示详细异常信息 +) + +# 导出配置好的logger +__all__ = ["logger"] diff --git a/models/__pycache__/log_alert.cpython-310.pyc b/models/__pycache__/log_alert.cpython-310.pyc index a77da78..119b1b8 100644 Binary files a/models/__pycache__/log_alert.cpython-310.pyc and b/models/__pycache__/log_alert.cpython-310.pyc differ diff --git a/models/log_alert.py b/models/log_alert.py index 791d8ad..548da62 100644 --- a/models/log_alert.py +++ b/models/log_alert.py @@ -10,12 +10,13 @@ class LogAlert(Base): __tablename__ = "log_alerts" id = Column(Integer, primary_key=True, index=True) - alter_name = Column(String(50)) - alter_timestamp = Column(Integer) + alter_name = Column(String(50), nullable=True) + alter_time = Column(DateTime, nullable=True) - timestamp = Column(Integer) - message = Column(Text) - status = Column(String(20), default="pending") # pending, notified, resolved + log_time = Column(DateTime, nullable=True) + log_topic = Column(String(50)) + log_message = Column(Text) + log_status = Column(String(20), default="pending") # pending, notified, resolved create_time = Column(DateTime, nullable=True, default=datetime.now) update_time = Column(DateTime, nullable=True, default=datetime.now, onupdate=datetime.now) diff --git a/requirements.txt b/requirements.txt index ad92be6..01b4c0f 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1 +1 @@ -fastapi~=0.119.0 pydantic~=2.12.2 SQLAlchemy~=2.0.44 requests~=2.32.5 pydantic-settings~=2.11.0 uvicorn~=0.38.0 \ No newline at end of file +fastapi~=0.119.0 pydantic~=2.12.2 SQLAlchemy~=2.0.44 requests~=2.32.5 pydantic-settings~=2.11.0 uvicorn~=0.38.0 loguru~=0.7.3 \ No newline at end of file diff --git a/schemas/__pycache__/log_alert.cpython-310.pyc b/schemas/__pycache__/log_alert.cpython-310.pyc index 2c2bf8a..90e1615 100644 Binary files a/schemas/__pycache__/log_alert.cpython-310.pyc and b/schemas/__pycache__/log_alert.cpython-310.pyc differ diff --git a/schemas/log_alert.py b/schemas/log_alert.py index b626702..c03d0a0 100644 --- a/schemas/log_alert.py +++ b/schemas/log_alert.py @@ -1,3 +1,5 @@ +from datetime import datetime + from pydantic import BaseModel, Field @@ -6,6 +8,21 @@ class LogAlertTrigger(BaseModel): timestamp: int +class LogAlertResponse(BaseModel): + alterName: str + alterTime: datetime + logTime: datetime + logTopic: str + logMessage: str + logStatus: str + + class Config: + from_attributes = True + json_encoders = { + # 自定义 datetime 类型的序列化格式 + datetime: lambda dt: dt.strftime('%Y-%m-%d %H:%M:%S') + } + class OpenobserveQuery(BaseModel): start_time: int end_time: int diff --git a/service/__pycache__/log_alert_service.cpython-310.pyc b/service/__pycache__/log_alert_service.cpython-310.pyc index e6c1bb9..dd5e30c 100644 Binary files a/service/__pycache__/log_alert_service.cpython-310.pyc and b/service/__pycache__/log_alert_service.cpython-310.pyc differ diff --git a/service/__pycache__/message_service.cpython-310.pyc b/service/__pycache__/message_service.cpython-310.pyc index 4fabce2..81956d0 100644 Binary files a/service/__pycache__/message_service.cpython-310.pyc and b/service/__pycache__/message_service.cpython-310.pyc differ diff --git a/service/log_alert_service.py b/service/log_alert_service.py index 1c91d37..50f492f 100644 --- a/service/log_alert_service.py +++ b/service/log_alert_service.py @@ -1,16 +1,19 @@ -from collections import defaultdict +from typing import List -from sqlalchemy import insert, select +from sqlalchemy import insert, select, func from sqlalchemy.orm import Session from models.log_alert import LogAlert -from schemas.log_alert import LogAlertTrigger +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") @@ -20,20 +23,19 @@ def trigger_alert(log_alter: LogAlertTrigger, db: Session) -> bool: return False # 判断日志告警是否已经存在 - log_db_alter = db.execute( - select(LogAlert) - .where(LogAlert.alter_timestamp == log_alter.timestamp) - ).first() + 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 - # 企业微信通知 - alter_message = f"🔔 **告警服务**:{log_alter.alter_name} \n 🕒 **告警时间**:{format_timestamp(log_alter.timestamp)}" - send_wechat_message(MessageEnum.MARKDOWN2, alter_message) - # 批量插入 - return batch_insert(log_alter, hits, db) + batch_insert(log_alter, hits, db) + + # 通知消息 + # notify_message(log_alter) + + return True def batch_insert(log_alter: LogAlertTrigger, hits: [], db: Session) -> bool: @@ -41,10 +43,11 @@ def batch_insert(log_alter: LogAlertTrigger, hits: [], db: Session) -> bool: [ { "alter_name": log_alter.alter_name, - "alter_timestamp": log_alter.timestamp, - "timestamp": hit["_timestamp"], - "message": hit["message"], - "status": "pending" + "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 ] @@ -56,15 +59,44 @@ def batch_insert(log_alter: LogAlertTrigger, hits: [], db: Session) -> bool: 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): - alerts = db.execute(select(LogAlert)).scalars().all() + 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()) + ) - # 按照name和timestamp分类 - classified = defaultdict(lambda: defaultdict(list)) - for alert in alerts: - classified[alert.alter_name][alert.alter_timestamp].append({ - "timestamp": format_timestamp(alert.timestamp), - "message": alert.message - }) + results = db.execute(stmt).fetchall() - return dict(classified) + 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] diff --git a/service/message_service.py b/service/message_service.py index c10d057..dd0772b 100644 --- a/service/message_service.py +++ b/service/message_service.py @@ -26,7 +26,8 @@ def send_wechat_message(message_type: MessageEnum, message: str): settings.wechat_webhook_url, headers={"Content-Type": "application/json"}, params={'key': settings.wechat_webhook_key}, - data=json.dumps(wechat_message, ensure_ascii=False).encode('utf-8') + data=json.dumps(wechat_message, ensure_ascii=False).encode('utf-8'), + timeout=5 ) return True diff --git a/utils/__pycache__/time.cpython-310.pyc b/utils/__pycache__/time.cpython-310.pyc index ef68244..5d20f36 100644 Binary files a/utils/__pycache__/time.cpython-310.pyc and b/utils/__pycache__/time.cpython-310.pyc differ diff --git a/utils/time.py b/utils/time.py index 0aac705..9aabd55 100644 --- a/utils/time.py +++ b/utils/time.py @@ -21,5 +21,5 @@ def get_timestamp_range(ts: int, delta_seconds: int = 5, unit: str = "microsecon return ts - delta, ts + delta -def format_timestamp(timestamp: int) -> str: - return datetime.fromtimestamp(timestamp / 1_000_000).strftime("%Y-%m-%d %H:%M:%S") +def format_timestamp(timestamp: int) -> datetime: + return datetime.fromtimestamp(timestamp / 1_000_000)